o
    ›¨Êhœ  ã                   @   sÀ   d 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fdd„Z dS )z'Task execution strategy (optimization).é    N)Úto_timestamp)Úsignals)Útrace)Ú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œ}| | jpoi ¡ |  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ÚupdateÚheaders)ÚmessageÚbodyr   r   r-   Úembedr   r   úH/var/www/html/env/lib/python3.10/site-packages/celery/worker/strategy.pyÚhybrid_to_proto2   sB   
ÿÿ

ñür2   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    r-   Útasksetr   r"   r#   r$   Nr%   Tr'   )r(   r)   r*   r   r+   r,   r   r-   )r.   r/   r   r   r0   r   r   r1   Úproto1_to_proto2B   s4   
ÿÿýÿür4   c	                    sØ   ˆj ‰ˆj‰t tj¡‰ˆj‰ˆoˆj}	ˆoˆj‰|	oˆj	‰ˆj
j‰ˆj‰ˆj ‰ˆjj‰	ˆj‰
ˆj‰ˆj‰tˆjƒ}
t|
ˆˆjˆˆˆd�‰ ˆ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.
    )Úappc                    sP  |d u rd| j vr| j| jdˆ ¡ f\}}}}nd| j v r(t| | j ƒ\}}}}n	ˆ| |ƒ\}}}}ˆ| ||ˆˆˆ	ˆˆ||||d�‰ ˆrZˆ jˆ jˆ jˆ jˆ j	dœ}	ˆt
j|	d|	id� ˆ jsbˆ jˆv rhˆ  ¡ rhd 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_rejectr5   ÚhostnameÚeventerr   Úconnection_errorsr/   r-   Údecodedr'   )r   Únamer   r   r   Údata)Úextra)ÚsenderÚrequestztask-receivedr   r   )	Úuuidr<   r   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   r1   Ú
<listcomp>Î   s    z9default.<locals>.task_message_handler.<locals>.<listcomp>)Úpayloadr/   r-   Úuses_utc_timezoner2   r   r<   r   r    r   Ú
_app_traceÚLOG_RECEIVEDr   Úrevokedr   Útask_receivedÚsendr   r   Úrequest_dictr(   Ú	isoformatr'   r	   ÚOverflowErrorÚ
ValueErrorÚinfoÚrejectÚqosÚincrement_eventually)r.   r/   ÚackrX   r"   r   r-   r;   r'   ÚcontextÚbucketr   Úexc©ÚReqÚ
_does_infor5   Úapply_eta_taskÚcall_atr:   ÚconsumerÚerrorr9   Ú
get_bucketÚhandler8   rW   Úlimit_post_etaÚ
limit_taskr4   Úrate_limits_enabledÚrevoked_tasksÚ
send_eventr   Útask_message_handlerr   Útask_sends_eventsÚto_system_tzrI   r1   rm   „   s†   ÿ
ÿüûù
€ÿ€ý

ÿ
z%default.<locals>.task_message_handler)r8   r:   ÚloggerÚisEnabledForÚloggingÚINFOÚevent_dispatcherÚenabledrR   Úsend_eventsÚtimerrc   rb   Údisable_rate_limitsÚtask_bucketsÚ__getitem__Úon_task_requestÚ_limit_taskÚ_limit_post_etar   ÚRequestr   ÚpoolÚ
controllerÚstaterP   r   )r   r5   rd   rW   re   r   ro   Úbytesr4   Úeventsr~   r   r_   r1   r   c   s(   





<ÿLr   )!Ú__doc__rr   Úkombu.asynchronous.timerr   Úceleryr   Ú
celery.appr   rN   Úcelery.exceptionsr   Úcelery.utils.importsr   Úcelery.utils.logr   Úcelery.utils.safereprr   Úcelery.utils.timer	   r@   r   r�   r   Ú__all__Ú__name__rp   r2   r4   rW   re   Ú	to_systemr‚   r   r   r   r   r1   Ú<module>   s(    )
"ý