o
    wvXjô  ã                   @   sÒ   d Z ddlmZm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 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dd„ Zdd„ Zejejeej e!eefdd„Z"dS )z'Task execution strategy (optimization).é    )Úabsolute_importÚunicode_literalsN)Úto_timestamp)Úbuffer_t)Úsignals)ÚInvalidTaskError)Úsymbol_by_name)Ú
get_logger)Úsaferepr)Útimezoneé   )Úcreate_request_cls)Útask_reserved)Údefaultc                 C   s  z|  dd¡|  di ¡}}|j W n ty   tdƒ‚ ty'   tdƒ‚w |  d¡|  d¡|  d¡|  d	¡|  d
¡|  d¡|  d¡|  d¡|  d¡|  d¡|  dd¡|  dd¡|  d¡|  d¡|  d¡dœ}|  d¡|  d¡|  d¡ddœ}|||f|d|  dd¡fS )zECreate a fresh protocol 2 message from a hybrid protocol 1/2 message.Úargs© Úkwargsú!Message does not have args/kwargsú(Task keyword arguments must be a mappingÚlangÚtaskÚidÚroot_idÚ	parent_idÚgroupÚmethÚshadowÚetaÚexpiresÚretriesr   Ú	timelimit)NNÚargsreprÚ
kwargsreprÚorigin)r   r   r   r   r   r   r   r   r   r   r   r    r!   r"   r#   Ú	callbacksÚerrbacksÚchordN©r$   r%   r&   ÚchainTÚutc)ÚgetÚitemsÚKeyErrorr   ÚAttributeError)ÚmessageÚbodyr   r   ÚheadersÚembedr   r   úS/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/worker/strategy.pyÚhybrid_to_proto2   s@   
ÿÿ

ñür3   c                 C   sÈ   z|  dd¡|  di ¡}}|j W n ty   tdƒ‚ ty'   tdƒ‚w |jt|ƒt|ƒ| jd� z|d |d< W n	 tyF   Y nw |  d	¡|  d
¡|  d¡ddœ}|||f|d|  dd¡fS )zŽConvert Task message protocol 1 arguments to protocol 2.

    Returns:
        Tuple: of ``(body, headers, already_decoded_status, utc)``
    r   r   r   r   r   )r!   r"   r0   Útasksetr   r$   r%   r&   Nr'   Tr)   )r*   r+   r,   r   r-   Úupdater
   r0   )r.   r/   r   r   r1   r   r   r2   Úproto1_to_proto2D   s4   
ÿÿýÿür6   c
                    sâ   ˆ	j ‰ˆ	j‰t tj¡‰ˆ	j‰ˆoˆj}
ˆoˆj‰|
oˆj	‰ˆ	j
j‰ˆ	j‰ˆ	j ‰ˆ	jj‰ˆ	j‰ˆ	j‰ˆ	j‰ˆ	jj‰tˆjƒ}t|ˆˆ	jˆˆƒ‰ ˆ	jjj‰tf‡ ‡‡‡‡‡‡‡‡‡	‡
‡‡‡‡‡‡‡‡‡‡‡‡‡‡‡‡fdd„	‰ˆS )z‘Default task execution strategy.

    Note:
        Strategies are here as an optimization, so sadly
        it's not very easy to override.
    c                    sF  |d u r%d| j vr%| j| jdˆ ¡ f\}}}}ˆs$t|ˆƒr"ˆ|ƒn|}nd| j v r5t| | j ƒ\}}}}n	ˆ| |ƒ\}}}}ˆ| ||ˆˆˆˆˆ	||||d�‰ ˆrUˆdˆ ƒ ˆ js]ˆ jˆv rcˆ  ¡ rcd S t	j
jˆ
ˆ d� ˆr’ˆdˆ jˆ jˆ jˆ jˆ jˆ jˆ j dd¡ˆ joˆˆ j ¡ ˆ jo�ˆ j ¡ d	�
 d }	d }
ˆ jrÖzˆ jr¥|ˆˆ jƒƒ}
n|ˆ jˆjƒ}
W n( ttfyÕ } zˆd
ˆ j|ˆ jdd�dd� ˆ jdd� W Y d }~nd }~ww ˆrÝˆˆjƒ}	|
rñ|	rñˆ
j ¡  ˆ|
ˆˆ |	dfdd�S |
�rˆ
j ¡  ˆ|
ˆˆ fdd� ˆS |	�rˆˆ |	dƒS ˆˆ ƒ |�r‡ fdd„|D ƒ ˆˆ ƒ d S )Nr   F)Úon_ackÚ	on_rejectÚappÚhostnameÚeventerr   Úconnection_errorsr/   r0   Údecodedr)   zReceived task: %s)ÚsenderÚrequestztask-receivedr   r   )	ÚuuidÚnamer   r   r   r   r   r   r   z2Couldn't convert ETA %r to timestamp: %r. Task: %rT)Úsafe)Úexc_info)Úrequeuer   é   )Úpriorityc                    s   g | ]}|ˆ ƒ‘qS r   r   )Ú.0Úcallback©Úreqr   r2   Ú
<listcomp>Ê   s    z9default.<locals>.task_message_handler.<locals>.<listcomp>)Úpayloadr/   r0   Úuses_utc_timezoneÚ
isinstancer3   r   r   Úrevokedr   Útask_receivedÚsendrA   r!   r"   r   r   Úrequest_dictr*   r   Ú	isoformatr)   r   ÚOverflowErrorÚ
ValueErrorÚinfoÚrejectÚqosÚincrement_eventually)r.   r/   ÚackrW   r$   r   r0   r=   r)   Úbucketr   Úexc©ÚReqÚ
_does_infor9   Úapply_eta_taskÚbody_can_be_bufferr   ÚbytesÚcall_atr<   ÚconsumerÚerrorr;   Ú
get_bucketÚhandler:   rV   Úlimit_post_etaÚ
limit_taskr6   Úrate_limits_enabledÚrevoked_tasksÚ
send_eventr   Útask_message_handlerr   Útask_sends_eventsÚto_system_tzrI   r2   rm   ‡   s€   ÿ€
ÿü
ù
€ÿ€ý

ÿ
z%default.<locals>.task_message_handler)r:   r<   ÚloggerÚisEnabledForÚloggingÚINFOÚevent_dispatcherÚenabledrQ   Úsend_eventsÚtimerrc   r`   Údisable_rate_limitsÚtask_bucketsÚ__getitem__Úon_task_requestÚ_limit_taskÚ_limit_post_etaÚpoolra   r   ÚRequestr   Ú
controllerÚstaterO   r   )r   r9   rd   rV   re   r   ro   rb   r   r6   Úeventsr   r   r]   r2   r   e   s*   





BÿEr   )#Ú__doc__Ú
__future__r   r   rr   Úkombu.asynchronous.timerr   Ú
kombu.fiver   Úceleryr   Úcelery.exceptionsr   Úcelery.utils.importsr   Úcelery.utils.logr	   Úcelery.utils.safereprr
   Úcelery.utils.timer   r?   r   r�   r   Ú__all__Ú__name__rp   r3   r6   rV   re   Ú	to_systemrb   r   r   r   r   r2   Ú<module>   s*   (
"ý