Ë
    ™V.j£q  ã                   óŽ  — 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#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l0m1Z1m2Z2 ddl3m4Z4 ddl5m6Z6m7Z7m8Z8m9Z9m:Z: dZ;ejx                  Z<ejz                  Z=e<e=hZ> e)e?«      Z@e@j‚                  e@j„                  e@j†                  e@jˆ                  e@jŠ                  f\  ZAZBZFZDZGdZHdZIdZJdZKdZLdZMdZNd ZOd!ZPd"ZQd#ZRd$„ ZS G d%„ d&«      ZT G d'„ d(ejª                  «      ZVy))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.
é    N)Údefaultdict)Úsleep)Úrestart_state)ÚRestartFreqExceeded)Ú	DummyLock)ÚContentDisallowedÚDecodeError)Ú_detect_environment)Ú	safe_repr)ÚTokenBucket)ÚppartialÚpromise)Ú	bootstepsÚsignals)Úbuild_tracer)ÚCPendingDeprecationWarningÚInvalidTaskErrorÚNotRegisteredÚWorkerShutdownÚWorkerTerminate)Únoop)Ú
get_logger)Úgethostname)ÚBunch)Útruncate)Úhumanize_secondsÚrate)Úloops)Úactive_requestsÚmaybe_shutdownÚrequestsÚ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
https://docs.celeryq.dev/en/latest/internals/protocol.html
for more information.

The full contents of the message body was:
%s

The 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:
https://docs.celeryq.dev/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                 ó’   — |€| j                   n|}dj                  t        t        |«      d«      t	        | j                   «      «      S )z+Format message body for debugging purposes.z{} ({}b)i   )ÚbodyÚformatr   r   Úlen)Úmr(   s     úaC:\xampp\htdocs\tradingbinance\backend\.venv\Lib\site-packages\celery/worker/consumer/consumer.pyr&   r&   ‚   s>   € ð �\ˆ1�6Š6 t€DØ×ÑœX¤i°£o°tÓ<Ü  §¡›[ó*ð *ó    c                   ó†  — e Zd ZdZeZdZdZdZdZ	dZ
 G d„ dej                  «      Zeddddddddddd	fd
„Zd„ Zd„ Zd„ Zd„ Zd3d„Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Z d„ Z!d„ Z"d„ Z#d „ Z$d4d!„Z%d4d"„Z&d#„ Z'd$„ Z(d%„ Z)	 	 d5d&„Z*d'„ Z+d(„ Z,d)„ Z-d*„ Z.d+„ Z/d,„ Z0d-„ Z1e2fd.„Z3d/„ Z4e5d0„ «       Z6e5d1„ «       Z7d2„ Z8y)6r$   úConsumer blueprint.NéÿÿÿÿTc                   ó"   — e Zd ZdZdZg d¢Zd„ Zy)ú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                 ó(   — | j                  |d«       y )NÚshutdown)Úsend_all)ÚselfÚparents     r,   r4   zConsumer.Blueprint.shutdown°   s   € Ø�M‰M˜& *Õ-r-   N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚnameÚdefault_stepsr4   © r-   r,   Ú	Blueprintr2       s   „ Ù!àˆò

ˆó	.r-   r?   Fé   é   c           	      óV  — || _         || _        || _        |xs
 t        «       | _        t        j                  «       | _        || _        || _	        | j                  «       | _        | j                   j                  «       | _        | j                  j                  | _        | j                  j                  | _        t!        dd¬«      | _        t$        j'                  t(        j*                  «      | _        d| _        || _        t3        «       | _        | j                   j6                  j8                  | _        || _        || _        || _         d| _!        tE        d„ «      | _#        | jI                  «        || _%        | jJ                  stM        | j                  dd«      r9|	| _'        | jN                  €-| j                   j6                  jP                  | _'        nd| _'        tS        | d	«      s'|rtT        jV                  ntT        jX                  | _-        t]        «       d
k(  rd | j                   j6                  _/        g | _0        g | _1        | je                  | j                   jb                  d   | jf                  ¬«      | _4         | jh                  jj                  | fi tm        |
xs i fi |¤Ž¤Ž y )Né   rA   )ÚmaxRÚmaxTr   Tc                   ó   — y ©Nr>   r>   r-   r,   Ú<lambda>z#Consumer.__init__.<locals>.<lambda>Ò   s   € °r-   Úis_greenFÚloopÚgeventÚconsumer)ÚstepsÚon_close)7Ú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_multiplierÚ_maximum_prefetch_restoredr   Útask_bucketsÚreset_rate_limitsÚhubÚgetattrÚamqheartbeatÚbroker_heartbeatÚhasattrr   ÚasynloopÚsynlooprJ   r
   Úbroker_connection_timeoutÚ_pending_operationsrM   r?   rN   Ú	blueprintÚapplyÚdict)r6   re   rQ   rR   rV   rO   rW   rP   rq   rs   Úworker_optionsrk   rl   rm   Úkwargss                  r,   Ú__init__zConsumer.__init__³   s  € ð ˆŒØ$ˆŒØ*ˆÔØ Ò1¤K£MˆŒÜ—9‘9“;ˆŒØˆŒ	ØˆŒ
ØŸ/™/Ó+ˆŒØŸ™×4Ñ4Ó6ˆŒØ!%§¡×!@Ñ!@ˆÔØ"Ÿm™m×:Ñ:ˆÔÜ+°¸Ô;ˆÔä ×-Ñ-¬g¯l©lÓ;ˆŒØˆÔØ.ˆÔÜ"›uˆÔØ!%§¡§¡×!IÑ!IˆÔØ#6ˆÔ Ø&<ˆÔ#Ø#6ˆÔ Ø*.ˆÔ'ô (©Ó5ˆÔØ×ÑÔ àˆŒØ�8Š8”w˜tŸy™y¨*°eÔ<Ø ,ˆDÔØ× Ñ Ð(Ø$(§H¡H§M¡M×$BÑ$B�Õ!à !ˆDÔä�t˜VÔ$Ù*-œŸš´5·=±=ˆDŒIäÓ  HÒ,ð 7;ˆD�H‰H�M‰MÔ3à#%ˆÔ àˆŒ
ØŸ™Ø—(‘(—.‘. Ñ,Ø—]‘]ð (ó 
ˆŒð 	ˆ�‰×Ñ˜TÑJ¤T¨.Ò*>¸BÑ%IÀ&Ñ%IÓJr-   c                 ó¨   — t        |g|¢­i |¤Ž}| j                  r| j                  j                  |«      S | j                  j	                  |«       |S rG   )r   rq   Ú	call_soonry   Úappend)r6   ÚpÚargsr~   s       r,   r�   zConsumer.call_soonï   sK   € Ü�QÐ(˜Ò( Ñ(ˆØ�8Š8Ø—8‘8×%Ñ% aÓ(Ð(Ø× Ñ ×'Ñ'¨Ô*Øˆr-   c                 óê   — | j                   s;| j                  r.	  | j                  j                  «       «        | j                  rŒ-y y y # t        $ r }t        j                  d|«       Y d }~Œ4d }~ww xY w)NzPending callback raised: %r)rq   ry   ÚpopÚ	Exceptionr_   Ú	exception©r6   Úexcs     r,   Úperform_pending_operationsz#Consumer.perform_pending_operationsö   sg   € Ø�xŠxØ×*Ò*ðIØ2�D×,Ñ,×0Ñ0Ó2Ô4ð ×*Õ*ð øô !ò IÜ×$Ñ$Ð%BÀC×HÑHûðIús   šA	 Á		A2ÁA-Á-A2c                 óP   — t        t        |dd «      «      }|rt        |d¬«      S d S )NÚ
rate_limitrA   )Úcapacity)r   rr   r   )r6   ÚtypeÚlimits      r,   Úbucket_for_taskzConsumer.bucket_for_taskþ   s)   € Ü”W˜T <°Ó6Ó7ˆÙ16Œ{˜5¨1Ô-Ð@¸DÐ@r-   c                 ó’   ‡ — ‰ j                   j                  ˆ fd„‰ j                  j                  j	                  «       D «       «       y )Nc              3   óJ   •K  — | ]  \  }}|‰j                  |«      f–— Œ y ­wrG   )r‘   )Ú.0ÚnÚtr6   s      €r,   Ú	<genexpr>z-Consumer.reset_rate_limits.<locals>.<genexpr>  s*   øè ø€ ð !
Ù5K©T¨Q°ˆQ�×$Ñ$ QÓ'Ô(Ñ5Kùs   ƒ #)ro   ÚupdaterO   ÚtasksÚitems©r6   s   `r,   rp   zConsumer.reset_rate_limits  s5   ø€ Ø×Ñ× Ñ ó !
Ø59·X±X·^±^×5IÑ5IÔ5Kó!
õ 	
r-   c                 ó¾   — | j                   j                  }| j                  r|sy| j                   j                  | j                  z  | _        | 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)rV   Únum_processesrl   rm   Ú_update_qos_eventually)r6   Úindexr�   s      r,   Ú_update_prefetch_countzConsumer._update_prefetch_count  sT   € ð Ÿ	™	×/Ñ/ˆØ×*Ò*±-Øà�I‰I×#Ñ# d×&>Ñ&>Ñ>ð 	Ô#ð ×*Ñ*¨5Ó1Ð1r-   c                 óœ   —  |dk  r| j                   j                  n| j                   j                  t        |«      | j                  z  «      S )Nr   )ÚqosÚdecrement_eventuallyÚincrement_eventuallyÚabsrm   )r6   rŸ   s     r,   rž   zConsumer._update_qos_eventually  sB   € ð3°¸²�—‘×-Ò-Ø—X‘X×2Ñ2Ü�‹J˜×1Ñ1Ñ1ó3ð 	3r-   c                 ó<   — t        |«       | j                  |«       y rG   )r#   re   )r6   Úrequests     r,   Ú_limit_move_to_poolzConsumer._limit_move_to_pool  s   € Ü�gÔØ×Ñ˜WÕ%r-   c                 ót  — 	 	 |j                  «       \  }}|j                  |«      r| j                  |«       Œ8|j                  j                  ||f«       | j                  dz   dz  x}| _        |j                  |«      }| j                  j                  || j                  |f|¬«       y # t        $ r Y y w xY w)NrA   é
   )Úpriority)r†   Ú
IndexErrorÚcan_consumer¨   ÚcontentsÚ
appendleftrd   Úexpected_timerW   Ú
call_afterÚ_schedule_bucket_request)r6   Úbucketr§   ÚtokensÚpriÚholds         r,   r²   z!Consumer._schedule_bucket_request#  sÂ   € ØðØ"(§*¡*£,‘�˜ð
 ×!Ñ! &Ô)Ø×(Ñ(¨Ô1Øð —‘×*Ñ*¨G°VÐ+<Ô=à+/×+<Ñ+<¸qÑ+@ÀBÑ*FÐF��dÔ'Ø×+Ñ+¨FÓ3�Ø—
‘
×%Ñ%Ø˜$×7Ñ7¸&¸Ø ð &ô ð
 øô% ò áðús   ƒB+ Â+	B7Â6B7c                 óJ   — |j                  ||f«       | j                  |«      S rG   )Úaddr²   ©r6   r§   r³   r´   s       r,   Ú_limit_taskzConsumer._limit_task;  s$   € Ø�
‰
�G˜VÐ$Ô%Ø×,Ñ,¨VÓ4Ð4r-   c                 ó~   — | j                   j                  «        |j                  ||f«       | j                  |«      S rG   )r¢   r£   r¸   r²   r¹   s       r,   Ú_limit_post_etazConsumer._limit_post_eta?  s4   € Ø�‰×%Ñ%Ô'Ø�
‰
�G˜VÐ$Ô%Ø×,Ñ,¨VÓ4Ð4r-   c                 óL  — | j                   }|j                  t        vr²t        «        | j                  r	 | j
                  j                  «        | xj                  dz  c_        | j                  j                  j                  r| j                  | j                  z   }n| j                  }	 |j                  | «       |j                  t        vrŒ±y y # t        $ r#}t        d|d¬«       t        d«       Y d }~Œ©d }~ww xY w# |$ �r}| j                   }d| _        | j#                  |«      }| j                  j                  |   }|s"t        d|rdnd› d|› d	�«       t%        d«      |‚t'        |t(        «      r4|j*                  t*        j,                  k(  rt        d
«       t/        d«      |‚t        «        |j                  t        vrP| j0                  r| j3                  |«       n| j5                  |«       | j7                  «        |j9                  | «       Y d }~�Œad }~ww xY w)NzFrequent restarts detected: %rrA   ©Úexc_infoFzRetrying to Ú	establishzre-establishzX a connection to the message broker after a connection loss has been disabled (app.conf.z=False). Shutting down...z Too many open files. Aborting...)rz   ÚstateÚSTOP_CONDITIONSr    Úrestart_countr^   Ústepr   Úcritr   rO   rh   Úbroker_channel_error_retryr\   r]   ÚstartÚfirst_connection_attemptÚ_get_connection_retry_typer   Ú
isinstanceÚOSErrorÚerrnoÚEMFILEr   Ú
connectionÚ#on_connection_error_after_connectedÚ$on_connection_error_before_connectedrN   Úrestart)r6   rz   rŠ   Úrecoverable_errorsÚis_connection_loss_on_startupÚconnection_retry_typeÚconnection_retrys          r,   rÇ   zConsumer.startD  sÑ  € Ø—N‘Nˆ	Ø�o‰o¤_Ñ4ÜÔØ×!Ò!ðØ×'Ñ'×,Ñ,Ô.ð ×Ò !Ñ#ÕØ�x‰x�}‰}×7Ò7Ø&*×&<Ñ&<¸t×?RÑ?RÑ&RÑ"à%)×%;Ñ%;Ð"ð,Ø—‘ Ô%ð �o‰o¤_Ô4øô
 +ò ÜÐ9¸3ÈÕKÜ˜!—H‘Hûðûð &ó ,ð 15×0MÑ0MÐ-Ø05�Ô-Ø(,×(GÑ(GÐHeÓ(fÐ%Ø#'§8¡8§=¡=Ð1FÑ#GÐ Ù'ÜØ&Ñ6S¡{ÐYgÐ&hð i3à3HÐ2IÐIbðdôô
 )¨Ó+°Ð4Ü˜c¤7Ô+°·	±	¼U¿\¹\Ò0IÜÐ;Ô<Ü)¨!Ó,°#Ð5ÜÔ Ø—?‘?¬/Ñ9Ø—’Ø×@Ñ@ÀÕEà×AÑAÀ#ÔFØ—M‘M”OØ×%Ñ% dÔ+ÿùð1,ús0   ¶C Â,D Ã	C>ÃC9Ã9C>ÄH#ÄDHÈH#c                 óN   — |r"| j                   j                  j                  �dS dS )NÚ"broker_connection_retry_on_startupÚbroker_connection_retry)rO   rh   r×   )r6   rÓ   s     r,   rÉ   z#Consumer._get_connection_retry_typeo  s-   € á1ØŸ™Ÿ™×HÑHÐTð 5ð 	0ð /ð	0r-   c                 óX   — t        t        | j                  j                  «       |d«       y )NzTrying to reconnect...)ÚerrorÚCONNECTION_ERRORr[   Úas_urir‰   s     r,   rÐ   z-Consumer.on_connection_error_before_connectedu  s!   € ÜÔ §¡× 4Ñ 4Ó 6¸Ø&õ	(r-   c           
      ó€  — t        t        d¬«       	 | j                  j                  «        | j
                  j                  j                  rdt        t        «      D ]Q  }|j                  j                  sŒ|j                  rŒ't        t        |«       |j                  | j                  «       ŒS nt!        j                   t"        t$        «       | j
                  j                  j&                  rÀt)        | j*                  | j,                  t/        t        t        «      «      | j*                  z  z
  «      | _        | j0                  | j,                  k(  | _        | j2                  sJt4        j7                  d| j0                  › dt/        t        t        «      «      › d| j,                  › d�«       y y y # t        $ r Y �Œ�w xY w)NTr¾   z+Temporarily reducing the prefetch count to z to avoid over-fetching since zW tasks are currently being processed.
The prefetch count will be gradually restored to z" as the tasks complete processing.)ÚwarnÚCONNECTION_RETRYrÎ   Úcollectr‡   rO   rh   Ú3worker_cancel_long_running_tasks_on_connection_lossÚtupler   ÚtaskÚ	acks_lateÚacknowledgedÚ3TERMINATING_TASK_ON_RESTART_AFTER_A_CONNECTION_LOSSÚcancelrV   ÚwarningsÚCANCEL_TASKS_BY_DEFAULTr   Ú&worker_enable_prefetch_count_reductionÚmaxrm   Úmax_prefetch_countr*   rl   rn   r_   Úinfo)r6   rŠ   r§   s      r,   rÏ   z,Consumer.on_connection_error_after_connectedy  sd  € ÜÔ¨Õ-ð	Ø�O‰O×#Ñ#Ô%ð �8‰8�=‰=×LÒLÜ ¤Ö1�Ø—<‘<×)Ó)°'×2FÓ2FÜÔLØ ô"à—N‘N 4§9¡9Õ-ñ	 2ô �M‰MÔ1Ô3MÔNà�8‰8�=‰=×?Ò?Ü*-Ø×(Ñ(Ø×'Ñ'¬#¬e´OÓ.DÓ*EÈ×H`ÑH`Ñ*`Ñ`ó+ˆDÔ'ð
 /3×.IÑ.IÈT×MdÑMdÑ.dˆDÔ+Ø×2Ò2Ü—‘ØAÀ$×B]ÑB]ÐA^ð _+Ü+.¬u´_Ó/EÓ+FÐ*Gð HHØHL×H_ÑH_ÐG`ð a+ð+õð 3ð @øô ò 	Úð	ús   “F0 Æ0	F=Æ<F=c                 óD   — | j                   j                  | d|fd¬«       y )NÚregister_with_event_loopzHub.register)r„   Údescription)rz   r5   )r6   rq   s     r,   rï   z!Consumer.register_with_event_loop˜  s&   € Ø�‰×ÑØÐ,°C°6Ø&ð 	 õ 	
r-   c                 ó:   — | j                   j                  | «       y rG   )rz   r4   r›   s    r,   r4   zConsumer.shutdownž  s   € Ø�‰×Ñ Õ%r-   c                 ó:   — | j                   j                  | «       y rG   )rz   Ústopr›   s    r,   ró   zConsumer.stop¡  s   € Ø�‰×Ñ˜DÕ!r-   c                 óB   — | j                   d c}| _         |r	 || «       y y rG   )rQ   )r6   Úcallbacks     r,   Úon_readyzConsumer.on_ready¤  s&   € Ø'+×'9Ñ'9¸4Ð$ˆ�$Ô$ÙÙ�T�Nð r-   c           	      óÌ   — | | j                   | j                  | j                  | j                  | j                  | j
                  | j                  j                  | j                  f	S rG   )	rÎ   Útask_consumerrz   rq   r¢   rs   rO   Úclockrj   r›   s    r,   Ú	loop_argszConsumer.loop_args©  sK   € Ø�d—o‘o t×'9Ñ'9Ø—‘ §¡¨$¯(©(°D×4EÑ4EØ—‘—‘ × 6Ñ 6ð8ð 	8r-   c                 óÆ   — t        t        ||j                  |j                  t	        |j
                  «      t        ||j                  «      d¬«       |j                  «        y)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.
        rA   r¾   N)	rÅ   ÚMESSAGE_DECODE_ERRORÚcontent_typeÚcontent_encodingr   Úheadersr&   r(   Úack)r6   ÚmessagerŠ   s      r,   Úon_decode_errorzConsumer.on_decode_error®  sI   € ô 	Ô!Ø�'×&Ñ&¨×(@Ñ(@Ü�w—‘Ó'¬°7¸G¿L¹LÓ)IØõ	ð 	�‰�r-   c                 ó  — | j                   r:| j                   j                  r$| j                   j                  j                  «        | j                  r| j                  j                  «        | j                  j                  «       D ]  }|sŒ|j                  «        Œ t        D ]  }|t        v sŒt        |= Œ t        j                  «        | j                  r2| j                  j                  r| j                  j                  «        y y y rG   )rP   Ú	semaphoreÚclearrW   ro   ÚvaluesÚclear_pendingr"   r!   rV   Úflush)r6   r³   Ú
request_ids      r,   rN   zConsumer.on_close¾  s¼   € ð �?Š?˜tŸ™×8Ò8Ø�O‰O×%Ñ%×+Ñ+Ô-Ø�:Š:Ø�J‰J×ÑÔØ×'Ñ'×.Ñ.Ö0ˆFÚØ×$Ñ$Õ&ð 1÷ ,ˆJØœXÒ%Ü˜ZÑ(ð ,ô 	×ÑÔ!Ø�9Š9˜Ÿ™ŸšØ�I‰I�O‰OÕð )ˆ9r-   c                 ó¶   — | j                  | j                  ¬«      }| j                  r0|j                  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)rZ   rs   rq   Ú	transportrï   rÎ   )r6   Úconns     r,   ÚconnectzConsumer.connectÐ  sE   € ð ×'Ñ'°$×2CÑ2CÐ'ÓDˆØ�8Š8Ø�N‰N×3Ñ3°D·O±OÀTÇXÁXÔNØˆr-   c                 óX   — | j                  | j                  j                  |¬«      «      S ©Nr  )Úensure_connectedrO   rZ   ©r6   r  s     r,   rZ   zConsumer.connection_for_readÛ  s*   € Ø×$Ñ$Ø�H‰H×(Ñ(°9Ð(Ó=ó?ð 	?r-   c                 óX   — | j                  | j                  j                  |¬«      «      S r  )r  rO   Úconnection_for_writer  s     r,   r  zConsumer.connection_for_writeß  s,   € Ø×$Ñ$Ø�H‰H×)Ñ)°IÐ)Ó>ó@ð 	@r-   c                 óx  ‡ ‡— t         fˆˆ fd„	}d}‰ j                  j                  j                  €b‰ j                  j                  j                   }t        j                  t        d‰ j                  j                  j                  › d�«      «       nO‰ j                  r"‰ j                  j                  j                   }n!‰ j                  j                  j                   }|r‰j                  «        d‰ _        ‰S ‰j                  |‰ j                  j                  j                  t        ¬«      Šd‰ _        ‰S )Nc                 ó  •— t        ‰dd «      r|dk(  rt        }|j                  t        |dd«      t	        |dz  «      ‰j
                  j                  j                  ¬«      }t        t        ‰j                  «       | |«       y )NÚaltr   ÚinÚ r@   )ÚwhenÚretriesÚmax_retries)rr   ÚCONNECTION_FAILOVERr)   r   ÚintrO   rh   Úbroker_connection_max_retriesrÚ   rÛ   rÜ   )rŠ   ÚintervalÚ	next_stepr  r6   s      €€r,   Ú_error_handlerz1Consumer.ensure_connected.<locals>._error_handleræ  sp   ø€ Ü�t˜U DÔ)¨h¸!ªmÜ/�	Ø!×(Ñ(Ü% h°°cÓ:Ü˜H q™LÓ)Ø ŸH™HŸM™M×GÑGð )ó IˆIô Ô" D§K¡K£M°3¸	ÕBr-   Fa$  The broker_connection_retry configuration setting will no longer determine
whether broker connection retries are made during startup in Celery 6.0 and above.
If you wish to retain the existing behavior for retrying connections on startup,
you should set broker_connection_retry_on_startup to Ú.)rõ   )ÚCONNECTION_RETRY_STEPrO   rh   r×   rØ   rè   rÞ   r   rÈ   r  Úensure_connectionr   r    )r6   r  r#  Úretry_disableds   ``  r,   r  zConsumer.ensure_connectedã  s  ù€ ô 5Jö 	Cð ˆà�8‰8�=‰=×;Ñ;ÐCð "&§¡§¡×!FÑ!FÐFˆNä�M‰MÜ*ðLð MQÏHÉHÏMÉM×LqÑLqÐKrÐrsðuóvõð ×,Ò,Ø%)§X¡X§]¡]×%UÑ%UÐ!U‘à%)§X¡X§]¡]×%JÑ%JÐ!J�áà�L‰LŒNØ,1ˆDÔ)ØˆKà×%Ñ%Ø˜DŸH™HŸM™M×GÑGÜ#ð &ó 
ˆð ).ˆÔ%Øˆr-   c                 óR   — | j                   r| j                   j                  «        y y rG   )Úevent_dispatcherr  r›   s    r,   Ú_flush_eventszConsumer._flush_events  s"   € Ø× Ò Ø×!Ñ!×'Ñ'Õ)ð !r-   c                 ó|   — | j                   r0| j                   j                  j                  | j                  «       y y rG   )rq   Ú_readyr¸   r*  r›   s    r,   Úon_send_event_bufferedzConsumer.on_send_event_buffered  s*   € Ø�8Š8Ø�H‰H�O‰O×Ñ × 2Ñ 2Õ3ð r-   c                 ó4  — | j                   }| j                  j                  j                  }||v r||   }n#|€|n|}|€dn|} |j                  |f|||dœ|¤Ž}|j                  |«      s.|j                  |«       |j                  «        t        d|«       y y )NÚdirect)ÚexchangeÚexchange_typeÚrouting_keyzStarted consuming from %s)	rø   rO   ÚamqpÚqueuesÚ
select_addÚconsuming_fromÚ	add_queueÚconsumerí   )	r6   Úqueuer0  r1  r2  ÚoptionsÚcsetr4  Úqs	            r,   Úadd_task_queuezConsumer.add_task_queue  s³   € à×!Ñ!ˆØ—‘—‘×%Ñ%ˆð �F‰?Ø�u‘‰Aà (Ð 0‘u°hˆHØ)6Ð)>™XØ"/ð à!�×!Ñ! %ð FØ+3Ø0=Ø.9ñFð >EñFˆAð ×"Ñ" 5Ô)Ø�N‰N˜1ÔØ�L‰LŒNÜÐ,¨eÕ4ð *r-   c                 ó°   — t        d|«       | j                  j                  j                  j	                  |«       | j
                  j                  |«       y )NzCanceling queue %s)rí   rO   r3  r4  Údeselectrø   Úcancel_by_queue)r6   r9  s     r,   Úcancel_task_queuezConsumer.cancel_task_queue4  s=   € ÜÐ! 5Ô)Ø�‰�‰×Ñ×%Ñ% eÔ,Ø×Ñ×*Ñ*¨5Õ1r-   c                 óp   — t        |«       | j                  |«       | j                  j                  «        y)zAMethod called by the timer to apply a task with an ETA/countdown.N)r#   re   r¢   r£   )r6   rã   s     r,   Úapply_eta_taskzConsumer.apply_eta_task9  s(   € ä�dÔØ×Ñ˜TÔ"Ø�‰×%Ñ%Õ'r-   c           	      óà   — t         j                  t        ||«      t        |j                  «      t        |j
                  «      t        |j                  «      t        |j                  «      «      S rG   )ÚMESSAGE_REPORTr)   r&   r   rý   rþ   Údelivery_inforÿ   ©r6   r(   r  s      r,   Ú_message_reportzConsumer._message_report?  sV   € Ü×$Ñ$¤Y¨w¸Ó%=Ü%.¨w×/CÑ/CÓ%DÜ%.¨w×/GÑ/GÓ%HÜ%.¨w×/DÑ/DÓ%EÜ%.¨w¯©Ó%?ó	Að 	Ar-   c                 óÈ   — t        t        | j                  ||«      «       |j                  t        | j
                  «       t        j                  j                  | |d ¬«       y )N©Úsenderr  rŠ   )	rÞ   ÚUNKNOWN_FORMATrH  Úreject_log_errorr_   r\   r   Útask_rejectedÚsendrG  s      r,   Úon_unknown_messagezConsumer.on_unknown_messageF  sJ   € ÜŒ^˜T×1Ñ1°$¸Ó@ÔAØ× Ñ ¤¨×)?Ñ)?Ô@Ü×Ñ×"Ñ"¨$¸ÀTÐ"ÕJr-   c           	      óú  — t        t        |t        ||«      |j                  |j                  d¬«       	 |j                  d   |j                  d   }}|j                  j                  d«      }t        |d ||j                  j                  d«      |j                  j                  d«      d ¬«      }|j                  t        | j                  «       | j                  j                  j                  |t!        |«      |¬	«       | j"                  r"| j"                  j%                  d
|d|›d�¬«       t&        j(                  j%                  | ||||¬«       y # t        $ r |j                  }|d   |d   }}d }Y �Œw xY w)NTr¾   Úidrã   Úroot_idÚcorrelation_idÚreply_to)r<   ÚchordrS  rT  rU  Úerrbacks)r§   ztask-failedzNotRegistered(Ú))Úuuidrˆ   )rK  r  rŠ   r<   rR  )rÚ   ÚUNKNOWN_TASK_ERRORr&   rÿ   rF  ÚgetÚKeyErrorÚpayloadr   Ú
propertiesrM  r_   r\   rO   ÚbackendÚmark_as_failurer   r)  rO  r   Útask_unknown)	r6   r(   r  rŠ   Úid_r<   rS  r]  r§   s	            r,   Úon_unknown_taskzConsumer.on_unknown_taskK  sl  € ÜÔ ØÜ˜ Ó&Ø�o‰oØ×#Ñ#Øõ	ð	ØŸ™¨Ñ-¨w¯©¸vÑ/F�ˆCØ—o‘o×)Ñ)¨)Ó4ˆGô
 Ø˜T¨7Ø"×-Ñ-×1Ñ1Ð2BÓCØ×'Ñ'×+Ñ+¨JÓ7Øô	
ˆð 	× Ñ ¤¨×)?Ñ)?Ô@Ø�‰×Ñ×(Ñ(Ø”˜tÓ$¨gð 	)ô 	
ð × Ò Ø×!Ñ!×&Ñ&Ø CØ*¨4¨(°!Ð4ð 'ô ô 	×Ñ×!Ñ!Ø ¨c¸Àð 	"õ 	
øô' ò 	Ø—o‘oˆGØ ™ w¨v¡�ˆCØ‹Gð	ús   µ9E Å!E:Å9E:c                 óÂ   — t        t        |t        ||«      d¬«       |j                  t        | j
                  «       t        j                  j                  | ||¬«       y )NTr¾   rJ  )	rÚ   ÚINVALID_TASK_ERRORr&   rM  r_   r\   r   rN  rO  )r6   r(   r  rŠ   s       r,   Úon_invalid_taskzConsumer.on_invalid_taskl  sL   € ÜÔ  #¤y°¸$Ó'?Øõ	à× Ñ ¤¨×)?Ñ)?Ô@Ü×Ñ×"Ñ"¨$¸ÀSÐ"ÕIr-   c                 ó,  — | j                   j                  }| j                   j                  j                  «       D ]W  \  }}|j	                  | j                   | «      | j
                  |<   t        |||| j                  | j                   ¬«      |_        ŒY y )N)rO   )	rO   Úloaderr™   rš   Ústart_strategyrY   r   rR   Ú	__trace__)r6   rh  r<   rã   s       r,   Úupdate_strategieszConsumer.update_strategiesr  sl   € Ø—‘—‘ˆØŸ(™(Ÿ.™.×.Ñ.Ö0‰JˆD�$Ø$(×$7Ñ$7¸¿¹À$Ó$GˆD�O‰O˜DÑ!Ü)¨$°°f¸d¿m¹mØ.2¯h©hô8ˆD�Nñ 1r-   c                 ó¾   ‡ ‡‡‡‡‡‡‡— ‰ j                   Š‰ j                  Š‰ j                  Š‰ j                  Š‰ j                  Š‰ j
                  Šˆˆˆˆˆˆˆ ˆfd„}|S )Nc                 ó¶  •— d }	 | j                   d   }	 ‰|   }	  ‰‰| j                  f‰j                  ¬«      } ‰‰| j                  f‰j                  ¬«      }‰j                  sv‰j                  dkD  rg‰j                  ‰j                  k  rN|j                  ‰j                  ‰j                  ¬«       |j                  ‰j                  ‰j                  ¬«        || |||‰«       y # t        $ r  ‰
d | «      cY S t        $ ri 	 | j                  «       }n*# t        $ r}‰j                  | |«      cY d }~cY S d }~ww xY w	 |d   |}}n # t        t        f$ r  ‰
|| «      cY cY S w xY wY �ŒZw xY w# t        t        f$ r} ‰	|| |«      cY d }~S d }~wt         $ r}‰j                  | |«      cY d }~S d }~ww xY w# t        $ r} ‰d | |«      cY d }~S d }~ww xY w)Nrã   )Úon_errorr   )rÿ   Ú	TypeErrorr\  Údecoder‡   r  Úack_log_errorÚ0_restore_prefetch_count_after_connection_restartrM  rn   rÃ   Ú_new_prefetch_countrì   Úthenr   r   r	   )r  r]  Útype_rŠ   ÚstrategyÚack_log_error_promiseÚreject_log_error_promiser�   Ú	callbacksrf  rP  rc  r   r6   rY   s          €€€€€€€€r,   Úon_task_receivedz6Consumer.create_task_handler.<locals>.on_task_received�  s  ø€ ð ˆGð@ØŸ™¨Ñ/�ð$>Ø% eÑ,�ð>Ù,3Ø!Ø ×.Ñ.Ð0Ø!%×!VÑ!Vô-Ð)ñ
 07Ø!Ø ×1Ñ1Ð3Ø!%×!VÑ!Vô0Ð,ð !×;Ò;Ø ×.Ñ.°Ò2Ø ×4Ñ4¸×8OÑ8OÒOà-×2Ñ2°4×3hÑ3hØ<@×<qÑ<qð 3ô sà0×5Ñ5°d×6kÑ6kØ?C×?tÑ?tð 6ô vñ Ø Ø-Ø0Ø!õ	øôM ò 9Ù)¨$°Ó8Ò8Üò @ð>Ø%Ÿn™nÓ.‘GøÜ ò >Ø×/Ñ/°¸Ó=×=ûð>úð@Ø%,¨V¡_°g˜7‘EøÜ!¤8Ð,ò @Ù-¨g°wÓ?Ô?ð@úò ð@ûôT )Ô*;Ð<ò BÙ*¨7°G¸SÓAÕAûÜ"ò >Ø×/Ñ/°¸Ó=Õ=ûð>ûôC ò ;Ù& t¨W°cÓ:Õ:ûð;ús¸   …C* •F; ›CE2 Ã*E/Ã>E/ÄDÄE/Ä	D?Ä!D:Ä2D?Ä3E/Ä:D?Ä?E/ÅEÅ
E/ÅE(Å#E/Å'E(Å(E/Å.E/Å2F8Æ	FÆ
F8ÆF8ÆF3Æ-F8Æ3F8Æ;	GÇ	GÇGÇG)rY   rP  rc  rf  rg   r�   )	r6   r   rz  r�   ry  rf  rP  rc  rY   s	   `` @@@@@@r,   Úcreate_task_handlerzConsumer.create_task_handlery  sU   ÿ€ Ø—_‘_ˆ
Ø!×4Ñ4ÐØ×.Ñ.ˆØ×.Ñ.ˆØ×(Ñ(ˆ	Ø—N‘Nˆ	÷5	>ó 5	>ðn  Ðr-   c                 óP  — | j                   j                  5  t        | j                  j                  j
                   | j                  f«      r
	 d d d «       y t        | j                  | j                  «      }|x| j                   _
        | _        | j                   j                  | j                   j                  «       | j                  }|| j                  k(  | _        |du r0| j                  du r"t        j                  d| j                  › �«       d d d «       y # 1 sw Y   y xY w)NFTzcResuming normal operations following a restart.
Prefetch count has been restored to the maximum of )r¢   Ú_mutexÚanyrO   rh   rê   rn   Úminrì   rs  Úvaluerl   rf   r_   rí   )r6   rƒ   r„   Únew_prefetch_countÚalready_restoreds        r,   rr  z9Consumer._restore_prefetch_count_after_connection_restartº  sì   € Ø�X‰X�_‹_ÜØ—H‘H—M‘M×HÑHÐHØ×/Ñ/ðô ð ÷ ˆ_ô "% T×%<Ñ%<¸d×>VÑ>VÓ!WÐØ;MÐMˆD�H‰HŒN˜TÔ8Ø�H‰H�L‰L˜Ÿ™Ÿ™Ô(à#×>Ñ>ÐØ.@ÀD×D[ÑD[Ñ.[ˆDÔ+à 5Ñ(¨T×-LÑ-LÐPTÑ-TÜ—‘ðJØJN×JaÑJaÐIbðdô÷ �_‰_ús   —8DÁB;DÄD%c                 óH   — | j                   j                  | j                  z  S rG   )rV   r�   rm   r›   s    r,   rì   zConsumer.max_prefetch_countÏ  s   € à�y‰y×&Ñ&¨×)AÑ)AÑAÐAr-   c                 óH   — | j                   j                  | j                  z   S rG   )r¢   r€  rm   r›   s    r,   rs  zConsumer._new_prefetch_countÓ  s   € à�x‰x�~‰~ × 8Ñ 8Ñ8Ð8r-   c                 óX   — dj                  | | j                  j                  «       ¬«      S )z``repr(self)``.z%<Consumer: {self.hostname} ({state})>)r6   rÁ   )r)   rz   Úhuman_stater›   s    r,   Ú__repr__zConsumer.__repr__×  s,   € à6×=Ñ=Ø˜TŸ^™^×7Ñ7Ó9ð >ó 
ð 	
r-   )r   rG   )NNN)9r8   r9   r:   r;   r|   rX   rQ   rV   rW   rÃ   rÈ   r   r?   r   r   r�   r‹   r‘   rp   r    rž   r¨   r²   rº   r¼   rÇ   rÉ   rÐ   rÏ   rï   r4   ró   rö   rú   r  rN   r  rZ   r  r  r*  r-  r=  rA  rC  rH  rP  rc  rf  rk  r   r{  rr  Úpropertyrì   rs  r‡  r>   r-   r,   r$   r$   Š   se  „ Ùà€Jð €Mð €Dð €Eà€Mð  $Ðô.�I×'Ñ'ô .ð(  $¨dØ Ø¨°$ÀTØ $¸%Ø()¸qó:KòxòIòAò
ó
2ò&3ò
&òò05ò5ò
),òV0ò(òò>
ò&ò"òò
8ò
ò ò$	ó?ó@ò1òf*ò4ð BFØ#'ó5ò,2ò
(òAòKò

òBJò8ð +2ó ? òBð* ñBó ðBð ñ9ó ð9ó
r-   r$   c                   ó$   — e Zd ZdZdZdZd„ Zd„ Zy)r%   zHEvent loop service.

    Note:
        This is always started last.
    z
event loopTc                 ó`   — | j                  |«        |j                  |j                  «       Ž  y rG   )Ú	patch_allrJ   rú   ©r6   Úcs     r,   rÇ   zEvloop.startè  s"   € Ø�‰�qÔØˆ�‰�—‘“Òr-   c                 ó6   — t        «       |j                  _        y rG   )r   r¢   r}  rŒ  s     r,   r‹  zEvloop.patch_allì  s   € Ü “{ˆ�‰�r-   N)r8   r9   r:   r;   ÚlabelÚlastrÇ   r‹  r>   r-   r,   r%   r%   Þ  s   „ ñð €EØ€Dòó#r-   r%   )Wr;   rÌ   ra   rS   rè   Úcollectionsr   Útimer   Úbilliard.commonr   Úbilliard.exceptionsr   Úkombu.asynchronous.semaphorer   Úkombu.exceptionsr   r	   Úkombu.utils.compatr
   Úkombu.utils.encodingr   Úkombu.utils.limitsr   Úviner   r   Úceleryr   r   Úcelery.app.tracer   Úcelery.exceptionsr   r   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!   r"   r#   Ú__all__ÚCLOSEÚ	TERMINATErÂ   r8   r_   Údebugrí   ÚwarningrÚ   ÚcriticalrÞ   rÅ   rß   r%  rÛ   r  rL  rZ  re  rü   rE  ræ   ré   r&   r$   ÚStartStopStepr%   r>   r-   r,   Ú<module>r­     s;  ðñó Û Û 	Û Ý #Ý å )Ý 3Ý 2ß ;Ý 2Ý *Ý *ß "ç %Ý )÷0õ 0å (Ý 'Ý .Ý &Ý &ß 4Ý ß kÕ kà
-€à�‰€Ø×Ñ€	Ø˜)Ð$€Ù	�HÓ	€Ø"(§,¡,°·±¸V¿^¹^Ø"(§,¡,°·±ð"AÑ €€tˆT�5˜$ðÐ ð
Ð ðÐ ð
Ð ð€ðÐ ð,
Ð ðÐ ð€ð7Ð 3ð
Ð ò*÷Q	
ñ Q	
ôh#ˆY×$Ñ$õ #r-   