o
    wvXjaX  ã                   @   sÆ  d Z ddlmZm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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mZmZmZmZmZm Z  ddl!m"Z"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- ddl.m/Z/ ddl0m1Z1m2Z2m3Z3 ddl4m5Z5 dZ6e7edƒZ8e+e9ƒZ:e:j;e:j<e:j=e:j>f\Z;Z<Z?Z>da@daAdd„ ZBeBƒ  e3jCZCejDjEZFe5jGZGe5jHZHe5jIZJe#G dd„ deKƒƒZLe	eJeHefdd„ZMdS )zfTask request.

This module defines the :class:`Request` class, that specifies
how tasks are executed.
é    )Úabsolute_importÚunicode_literalsN)Údatetime)Útime)Úref)ÚTERM_SIGNAME)Ú	safe_reprÚsafe_str)Úcached_property)Úsignals)ÚContext)Ú
trace_taskÚtrace_task_ret)ÚIgnoreÚInvalidTaskErrorÚRejectÚRetryÚTaskRevokedErrorÚ
TerminatedÚTimeLimitExceededÚWorkerLostError)Ú	monotonicÚpython_2_unicode_compatibleÚstring)ÚmaybeÚnoop)Ú
get_logger)Úgethostname)Úget_pickled_exception)Úmaybe_iso8601Úmaybe_make_awareÚtimezoneé   )Ústate)ÚRequestÚpypy_version_infoFc                   C   s   t  tj¡at  tj¡ad S ©N)ÚloggerÚisEnabledForÚloggingÚDEBUGÚ_does_debugÚINFOÚ
_does_info© r.   r.   úR/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/worker/request.pyÚ__optimize__1   s   r0   c                   @   sÖ  e Zd ZdZdZdZdZdZdZdZ	dZ
dZesdZeddddddeddddeefdd„Zed	d
„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd „ ƒZed!d"„ ƒZed#d$„ ƒZ ed%d&„ ƒZ!ed'd(„ ƒZ"ed)d*„ ƒZ#e#j$d+d*„ ƒZ#ed,d-„ ƒZ%ed.d/„ ƒZ&e&j$d0d/„ ƒZ&ed1d2„ ƒZ'ed3d4„ ƒZ(ed5d6„ ƒZ)ed7d8„ ƒZ*e*j$d9d8„ ƒZ*ed:d;„ ƒZ+ed<d=„ ƒZ,ed>d?„ ƒZ-e-j$d@d?„ ƒZ-edAdB„ ƒZ.e.j$dCdB„ ƒZ.edDdE„ ƒZ/edFdG„ ƒZ0dHdI„ Z1ddJdK„Z2dLdM„ Z3dvdNdO„Z4dPdQ„ Z5dRdS„ Z6dTdU„ Z7dVdW„ Z8dXdY„ Z9dZd[„ Z:d\d]„ Z;dwd^d_„Z<d`da„ Z=dxdbdc„Z>dxddde„Z?dfdg„ Z@dhdi„ ZAdjdk„ ZBeCdldm„ ƒZDeCdndo„ ƒZEeCdpdq„ ƒZFeCdrds„ ƒZGeCdtdu„ ƒZHdS )yr$   zA request for task execution.FN)NN)Ú_appÚ_typeÚnameÚidÚ_root_idÚ
_parent_idÚ_on_ackÚ_bodyÚ	_hostnameÚ_eventerÚ_connection_errorsÚ_taskÚ_etaÚ_expiresÚ_request_dictÚ
_on_rejectÚ_utcÚ_content_typeÚ_content_encodingÚ	_argsreprÚ_kwargsreprÚ_argsÚ_kwargsÚ_decodedÚ	__payloadÚ__weakref__Ú__dict__Tc              
   K   sª  || _ |d u r
|jn|| _|
d u r|jn|
| _|| _|| _|| _|r)d  | _| _	n	|j
|j| _| _	| jr8| jn|j| _| jd | _| jd  | _| _d| jv rY| jd pW| j| _| j d¡| _| j d¡| _| j dd ¡}|rs|| _| j dd¡| _| j d	d¡| _|| _|	| _|p�tƒ | _|| _|p•d
| _|pŸ| jj| j | _| j d¡}|d urÑz||ƒ}W n tt t!fyÈ } zt"d #||¡ƒ‚d }~ww ||| j$ƒ| _%nd | _%| j d¡}|d u�rz||ƒ}W n tt t!fyü } zt"d #||¡ƒ‚d }~ww ||| j$ƒ| _&nd | _&|j'�pi }|j(�pi }| d¡| d¡| d¡| d¡dœ| _)| j *| d¡| d¡| j| j)dœ¡ | j\| jd< | jd< }| jd | _+| jd | _,d S )Nr4   ÚtaskÚshadowÚroot_idÚ	parent_idÚ	timelimitÚargsreprÚ Ú
kwargsreprr.   Úetazinvalid ETA value {0!r}: {1}Úexpiresz invalid expires value {0!r}: {1}ÚexchangeÚrouting_keyÚpriorityÚredelivered)rV   rW   rX   rY   Úreply_toÚcorrelation_id)rZ   r[   ÚhostnameÚdelivery_infoÚargsÚkwargs)-Ú_messageÚheadersr?   Úbodyr8   r1   rA   rH   rB   rC   Úcontent_typeÚcontent_encodingÚpayloadÚ_Request__payloadr4   r2   r3   Úgetr5   r6   Útime_limitsrD   rE   r7   r@   r   r9   r:   r;   Útasksr<   ÚAttributeErrorÚ
ValueErrorÚ	TypeErrorr   ÚformatÚtzlocalr=   r>   r]   Ú
propertiesÚ_delivery_infoÚupdaterF   rG   )ÚselfÚmessageÚon_ackr\   ÚeventerÚappÚconnection_errorsÚrequest_dictrL   Ú	on_rejectrb   ra   ÚdecodedÚutcr    r   ÚoptsrP   rT   ÚexcrU   r]   ro   Ú_r.   r.   r/   Ú__init__[   sˆ   
ÿ


ÿ€ÿ

ÿ€ÿüüzRequest.__init__c                 C   ó   | j S r&   )rp   ©rr   r.   r.   r/   r]   ¬   ó   zRequest.delivery_infoc                 C   r€   r&   )r`   r�   r.   r.   r/   rs   °   r‚   zRequest.messagec                 C   r€   r&   ©r?   r�   r.   r.   r/   rx   ´   r‚   zRequest.request_dictc                 C   r€   r&   )r8   r�   r.   r.   r/   rb   ¸   r‚   zRequest.bodyc                 C   r€   r&   )r1   r�   r.   r.   r/   rv   ¼   r‚   zRequest.appc                 C   r€   r&   )rA   r�   r.   r.   r/   r{   À   r‚   zRequest.utcc                 C   r€   r&   )rB   r�   r.   r.   r/   rc   Ä   r‚   zRequest.content_typec                 C   r€   r&   )rC   r�   r.   r.   r/   rd   È   r‚   zRequest.content_encodingc                 C   r€   r&   )r2   r�   r.   r.   r/   ÚtypeÌ   r‚   zRequest.typec                 C   r€   r&   )r5   r�   r.   r.   r/   rN   Ð   r‚   zRequest.root_idc                 C   r€   r&   )r6   r�   r.   r.   r/   rO   Ô   r‚   zRequest.parent_idc                 C   r€   r&   )rD   r�   r.   r.   r/   rQ   Ø   r‚   zRequest.argsreprc                 C   r€   r&   )rF   r�   r.   r.   r/   r^   Ü   r‚   zRequest.argsc                 C   r€   r&   )rG   r�   r.   r.   r/   r_   à   r‚   zRequest.kwargsc                 C   r€   r&   )rE   r�   r.   r.   r/   rS   ä   r‚   zRequest.kwargsreprc                 C   r€   r&   )r7   r�   r.   r.   r/   rt   è   r‚   zRequest.on_ackc                 C   r€   r&   ©r@   r�   r.   r.   r/   ry   ì   r‚   zRequest.on_rejectc                 C   ó
   || _ d S r&   r…   ©rr   Úvaluer.   r.   r/   ry   ð   ó   
c                 C   r€   r&   )r9   r�   r.   r.   r/   r\   ô   r‚   zRequest.hostnamec                 C   r€   r&   ©r:   r�   r.   r.   r/   ru   ø   r‚   zRequest.eventerc                 C   r†   r&   rŠ   )rr   ru   r.   r.   r/   ru   ü   r‰   c                 C   r€   r&   )r;   r�   r.   r.   r/   rw      r‚   zRequest.connection_errorsc                 C   r€   r&   )r<   r�   r.   r.   r/   rL     r‚   zRequest.taskc                 C   r€   r&   )r=   r�   r.   r.   r/   rT     r‚   zRequest.etac                 C   r€   r&   ©r>   r�   r.   r.   r/   rU     r‚   zRequest.expiresc                 C   r†   r&   r‹   r‡   r.   r.   r/   rU     r‰   c                 C   s   | j d u r| jjj| _ | j S r&   )Ú_tzlocalr1   Úconfr!   r�   r.   r.   r/   rn     s   
zRequest.tzlocalc                 C   s   | j j p| j jS r&   )rL   Úignore_resultÚstore_errors_even_if_ignoredr�   r.   r.   r/   Ústore_errors  s   
ÿzRequest.store_errorsc                 C   r€   r&   ©r4   r�   r.   r.   r/   Útask_id  ó   zRequest.task_idc                 C   r†   r&   r‘   r‡   r.   r.   r/   r’   $  r‰   c                 C   r€   r&   ©r3   r�   r.   r.   r/   Ú	task_name(  r“   zRequest.task_namec                 C   r†   r&   r”   r‡   r.   r.   r/   r•   -  r‰   c                 C   ó
   | j d S )NrZ   rƒ   r�   r.   r.   r/   rZ   1  ó   
zRequest.reply_toc                 C   r–   )Nr[   rƒ   r�   r.   r.   r/   r[   6  r—   zRequest.correlation_idc                 K   s|   | j }| j}|  ¡ rt|ƒ‚| j\}}|jt| j|| j| j	| j
| jf| j| j| j| j|p.|j|p2|j|d�	}tt|ƒ| _|S )a  Used by the worker to send this task to the pool.

        Arguments:
            pool (~celery.concurrency.base.TaskPool): The execution pool
                used to execute this request.

        Raises:
            celery.exceptions.TaskRevokedError: if the task was revoked.
        ©r^   Úaccept_callbackÚtimeout_callbackÚcallbackÚerror_callbackÚsoft_timeoutÚtimeoutr[   )r4   r<   Úrevokedr   rh   Úapply_asyncr   r2   r?   r8   rB   rC   Úon_acceptedÚ
on_timeoutÚ
on_successÚ
on_failureÚsoft_time_limitÚ
time_limitr   r   Ú_apply_result)rr   Úpoolr_   r’   rL   r¦   r¥   Úresultr.   r.   r/   Úexecute_using_pool;  s(   

ÿözRequest.execute_using_poolc              
   C   s„   |   ¡ rdS | jjs|  ¡  | j\}}}| j}|j||ddœfi |p#i ¤Ž t| j| j| j	| j
|| j| jj| jd�d }|  ¡  |S )zÌExecute the task in a :func:`~celery.app.trace.trace_task`.

        Arguments:
            loglevel (int): The loglevel used by the task.
            logfile (str): The logfile used by the task.
        NF)ÚloglevelÚlogfileÚis_eager)r\   Úloaderrv   r   )rŸ   rL   Ú	acks_lateÚacknowledgeÚ_payloadr?   rq   r   r4   rF   rG   r9   r1   r®   )rr   r«   r¬   r~   ÚembedÚrequestÚretvalr.   r.   r/   Úexecute[  s*   ýü
þþzRequest.executec                 C   s6   | j rt | j j¡}|| j krt | j¡ dS dS dS )z%If expired, mark the task as revoked.TN)r>   r   ÚnowÚtzinfoÚrevoked_tasksÚaddr4   )rr   r¶   r.   r.   r/   Úmaybe_expirex  s   
üzRequest.maybe_expirec                 C   sn   t  |pt¡}| jr| | j|¡ |  dd|d¡ n||f| _| jd ur3|  ¡ }|d ur5| 	|¡ d S d S d S )NÚ
terminatedTF)
Ú_signalsÚsignumr   Ú
time_startÚterminate_jobÚ
worker_pidÚ_announce_revokedÚ_terminate_on_ackr§   Ú	terminate)rr   r¨   ÚsignalÚobjr.   r.   r/   rÃ   €  s   

ýzRequest.terminatec                 C   s^   t | ƒ | jd|||d� | jjj| j|| j| jd� |  ¡  d| _	t
| j| j|||d� d S )Nztask-revoked)r»   r½   Úexpired©r³   Ústore_resultT)r³   r»   r½   rÆ   )Ú
task_readyÚ
send_eventrL   ÚbackendÚmark_as_revokedr4   Ú_contextr�   r°   Ú_already_revokedÚsend_revoked)rr   Úreasonr»   r½   rÆ   r.   r.   r/   rÁ   Œ  s   ÿ
þ

ÿzRequest._announce_revokedc                 C   sV   d}| j rdS | jr|  ¡ }| jtv r)td| j| jƒ |  |r!dnddd|¡ dS dS )z%If revoked, skip task and mark state.FTzDiscarding revoked task: %s[%s]rÆ   rŸ   N)rÎ   r>   rº   r4   r¸   Úinfor3   rÁ   )rr   rÆ   r.   r.   r/   rŸ   ™  s   
ÿzRequest.revokedc                 K   s@   | j r| j jr| jjr| j j|fd| ji|¤Ž d S d S d S d S )NÚuuid)r:   ÚenabledrL   Úsend_eventsÚsendr4   )rr   r„   Úfieldsr.   r.   r/   rÊ   ¨  s   ÿzRequest.send_eventc                 C   sn   || _ tƒ tƒ |  | _t| ƒ | jjs|  ¡  |  d¡ t	r(t
d| j| j|ƒ | jdur5| j| jŽ  dS dS )z4Handler called when task is accepted by worker pool.ztask-startedzTask accepted: %s[%s] pid:%rN)rÀ   r   r   r¾   Útask_acceptedrL   r¯   r°   rÊ   r+   Údebugr3   r4   rÂ   rÃ   )rr   ÚpidÚtime_acceptedr.   r.   r/   r¡   ¬  s   

ÿzRequest.on_acceptedc                 C   s|   |rt d|| j| jƒ dS t| ƒ td|| j| jƒ t|ƒ}| jjj| j|| j	| j
d� | jjr:| jjr<|  ¡  dS dS dS )z%Handler called if the task times out.z)Soft time limit (%ss) exceeded for %s[%s]z)Hard time limit (%ss) exceeded for %s[%s]rÇ   N)Úwarnr3   r4   rÉ   Úerrorr   rL   rË   Úmark_as_failurerÍ   r�   r¯   Úacks_on_failure_or_timeoutr°   )rr   Úsoftrž   r}   r.   r.   r/   r¢   º  s    
ÿ
ÿ
þÿzRequest.on_timeoutc                 K   s^   |\}}}|rt |jttfƒr|j‚| j|dd�S t| ƒ | jjr%|  ¡  | j	d||d� dS )z6Handler called if the task was successfully processed.T©Ú	return_okútask-succeeded©r©   ÚruntimeN)
Ú
isinstanceÚ	exceptionÚ
SystemExitÚKeyboardInterruptr¤   rÉ   rL   r¯   r°   rÊ   ©rr   Úfailed__retval__runtimer_   Úfailedr´   rä   r.   r.   r/   r£   Í  s   
zRequest.on_successc                 C   s2   | j jr|  ¡  | jdt|jjƒt|jƒd� dS )z-Handler called if the task should be retried.ztask-retried©ræ   Ú	tracebackN)	rL   r¯   r°   rÊ   r   ræ   r}   r	   rí   )rr   Úexc_infor.   r.   r/   Úon_retryÛ  s   

þzRequest.on_retryc                 C   sV  t | ƒ t|jtƒrtd|jf ƒ‚t|jtƒr | j|jjd�S t|jtƒr*|  ¡ S |j}t|t	ƒr7|  
|¡S d}| jjrd| jjoEt|tƒ}| jj}|rWd}| j|d� d}n|r^|  ¡  n| jdd� t|tƒrv|  ddt|ƒd¡ d}n|s�t|tƒs|s�| jjj| j|| j| jd� |r�| jdtt|jƒƒ|jd� |s©td	||jd
� dS dS )z/Handler called if the task raised an exception.zProcess got: %s©ÚrequeueFTr»   rÇ   ztask-failedrì   zTask handler raised error: %r)rî   N)rÉ   rå   ræ   ÚMemoryErrorr   Úrejectrñ   r   r°   r   rï   rL   r¯   Úreject_on_worker_lostr   rÞ   r   rÁ   r   rË   rÝ   r4   rÍ   r�   rÊ   r   r   rí   rÜ   rî   )rr   rî   Úsend_failed_eventrá   r}   rñ   ró   Úackr.   r.   r/   r¤   ä  sX   

þ

ÿ
þý
ÿÿzRequest.on_failurec                 C   s"   | j s|  t| j¡ d| _ dS dS )zAcknowledge task.TN)Úacknowledgedr7   r'   r;   r�   r.   r.   r/   r°     s   
þzRequest.acknowledgec                 C   s2   | j s|  t| j|¡ d| _ | jd|d� d S d S )NTztask-rejectedrð   )r÷   r@   r'   r;   rÊ   )rr   rñ   r.   r.   r/   ró   $  s
   ýzRequest.rejectc                 C   s.   | j | j| j| j| j| j| j| j| j| j	dœ
S )N)
r4   r3   r^   r_   r„   r\   r¾   r÷   r]   rÀ   )
r4   r3   rF   rG   r2   r9   r¾   r÷   r]   rÀ   )rr   Úsafer.   r.   r/   rÑ   *  s   özRequest.infoc                 C   s
   d  | ¡S )Nz{0.name}[{0.id}])rm   r�   r.   r.   r/   Ú	humaninfo8  s   
zRequest.humaninfoc                 C   s<   d  |  ¡ | jrd | j¡nd| jrd | j¡g¡S dg¡S )z``str(self)``.ú z
 ETA:[{0}]rR   z expires:[{0}])Újoinrù   r=   rm   r>   r�   r.   r.   r/   Ú__str__;  s   ýýzRequest.__str__c                 C   s   d  t| ƒj|  ¡ | j| j¡S )z``repr(self)``.z<{0}: {1} {2} {3}>)rm   r„   Ú__name__rù   rD   rE   r�   r.   r.   r/   Ú__repr__C  s   þzRequest.__repr__c                 C   r€   r&   )rf   r�   r.   r.   r/   r±   J  r‚   zRequest._payloadc                 C   ó   | j \}}}| d¡S )NÚchord©r±   rg   ©rr   r~   r²   r.   r.   r/   r   N  ó   
zRequest.chordc                 C   rÿ   )NÚerrbacksr  r  r.   r.   r/   r  W  r  zRequest.errbacksc                 C   s   | j  d¡S )NÚgroup)r?   rg   r�   r.   r.   r/   r  `  s   zRequest.groupc                 C   s.   | j }| j\}}}|jdi |pi ¤Ž t|ƒS )z9Context (:class:`~celery.app.task.Context`) of this task.Nr.   )r?   r±   rq   r   )rr   r³   r~   r²   r.   r.   r/   rÍ   f  s   zRequest._contextr&   )TF)F)Irý   Ú
__module__Ú__qualname__Ú__doc__r÷   r¾   rÀ   rh   rÎ   rÂ   r§   rŒ   ÚIS_PYPYÚ	__slots__r   r    r   r   Úpropertyr]   rs   rx   rb   rv   r{   rc   rd   r„   rN   rO   rQ   r^   r_   rS   rt   ry   Úsetterr\   ru   rw   rL   rT   rU   rn   r�   r’   r•   rZ   r[   rª   rµ   rº   rÃ   rÁ   rŸ   rÊ   r¡   r¢   r£   rï   r¤   r°   ró   rÑ   rù   rü   rþ   r
   r±   r   r  r  rÍ   r.   r.   r.   r/   r$   D   sè    	
úQ


































 

	:





r$   c	           
   
      sJ   |j ‰|j‰|j‰|j‰ |o|j‰G ‡ ‡‡‡‡‡‡‡‡f	dd„d| ƒ}	|	S )Nc                       s2   e Zd Z‡‡‡‡‡‡fdd„Z‡ ‡‡fdd„ZdS )z#create_request_cls.<locals>.Requestc                    s~   | j }| js
|ˆv r|  ¡ rt|ƒ‚| j\}}ˆ ˆ| j|| j| j| j| j	f| j
| j| j| j|p0ˆ|p3ˆ|d�	}tˆ|ƒ| _|S )Nr˜   )r’   rU   rŸ   r   rh   r„   rx   rb   rc   rd   r¡   r¢   r£   r¤   r   r§   )rr   r¨   r_   r’   r¦   r¥   r©   )r    Údefault_soft_time_limitÚdefault_time_limitr   r¸   Útracer.   r/   rª   |  s&   
ÿöz6create_request_cls.<locals>.Request.execute_using_poolc                    sb   |\}}}|rt |jttfƒr|j‚| j|dd�S ˆ| ƒ ˆ r#|  ¡  ˆr/| jd||d� d S d S )NTrà   râ   rã   )rå   ræ   rç   rè   r¤   r°   rÊ   ré   )r¯   ÚeventsrÉ   r.   r/   r£   “  s   
ÿ
ÿÿz.create_request_cls.<locals>.Request.on_successN)rý   r  r  rª   r£   r.   ©	r¯   r    r  r  r  r   r¸   rÉ   r  r.   r/   r$   z  s    r$   )r¦   r¥   r    r¯   rÓ   )
ÚbaserL   r¨   r\   ru   r   r¸   rÉ   r  r$   r.   r  r/   Úcreate_request_clsq  s   
$*r  )Nr  Ú
__future__r   r   r)   Úsysr   r   Úweakrefr   Úbilliard.commonr   Úkombu.utils.encodingr   r	   Úkombu.utils.objectsr
   Úceleryr   Úcelery.app.taskr   Úcelery.app.tracer   r   Úcelery.exceptionsr   r   r   r   r   r   r   r   Úcelery.fiver   r   r   Úcelery.platformsr¼   Úcelery.utils.functionalr   r   Úcelery.utils.logr   Úcelery.utils.nodenamesr   Úcelery.utils.serializationr   Úcelery.utils.timer   r    r!   rR   r#   Ú__all__Úhasattrr	  rý   r'   rØ   rÑ   ÚwarningrÜ   rÛ   r-   r+   r0   Útz_or_localÚtask_revokedrÕ   rÏ   r×   rÉ   rŸ   r¸   Úobjectr$   r  r.   r.   r.   r/   Ú<module>   s\   (
ÿ    1þ