o
    ›¨Êh˜  ã                   @   sv  d Z ddlZddlZddlZddlmZ ddlmZmZm	Z	m
Z
 ddl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 ej d	d
¡Zedi d�Zedddhd�Zeddhd�ZG dd„ dejƒZeddeddddfdd„ƒZeddededdfdedede de
e ef de	e  dede!d e"d!eej fd"d#„ƒZ#eddedfd$d%„ƒZ$dede
e ef de d!dfd&d'„Z%dS )(z'Embedded workers for integration tests.é    N)Úcontextmanager)ÚAnyÚIterableÚOptionalÚUnion)ÚCeleryÚworker)Ú_set_task_join_will_blockÚallow_join_result)ÚSignal)Úanon_nodenameÚWORKER_LOGLEVELÚerrorÚtest_worker_starting)ÚnameÚproviding_argsÚtest_worker_startedr   ÚconsumerÚtest_worker_stoppedc                       sT   e Zd ZdZdZ‡ fdd„ZG dd„ dejjƒZ‡ fdd„Z	d	d
„ Z
dd„ Z‡  ZS )ÚTestWorkControllerz3Worker that can synchronize on being fully started.Nc                    s¤   t  ¡ | _tƒ j|i |¤Ž | jj d¡d dkrPddlm	} |ƒ | _
t ¡ | _zddlm} | ¡  W n	 ty=   Y nw tj | j
t ¡ ¡| _| j ¡  d S d S )NÚ.éÿÿÿÿÚpreforkr   )ÚQueue)Úpickling_support)Ú	threadingÚEventÚ_on_startedÚsuperÚ__init__Úpool_clsÚ
__module__ÚsplitÚbilliardr   Úlogger_queueÚosÚgetpidÚpidÚtblibr   ÚinstallÚImportErrorÚloggingÚhandlersÚQueueListenerÚ	getLoggerÚqueue_listenerÚstart)ÚselfÚargsÚkwargsr   r   ©Ú	__class__© úO/var/www/html/env/lib/python3.10/site-packages/celery/contrib/testing/worker.pyr   #   s   

ÿòzTestWorkController.__init__c                   @   s   e Zd Zdd„ Zdd„ ZdS )zTestWorkController.QueueHandlerc                 C   s
   d|_ |S )NT)Ú
from_queue©r1   Úrecordr6   r6   r7   Úprepare:   s   z'TestWorkController.QueueHandler.preparec                 C   s   t jr‚ d S )N)r+   ÚraiseExceptionsr9   r6   r6   r7   ÚhandleError?   s   ÿz+TestWorkController.QueueHandler.handleErrorN)Ú__name__r!   Ú__qualname__r;   r=   r6   r6   r6   r7   ÚQueueHandler9   s    r@   c                    s@   ˆ j rˆ  ˆ j ¡}| ‡ fdd„¡ t ¡ }| |¡ tƒ  ¡ S )Nc                    s   | j ˆ jkot| ddƒ S )Nr8   F)Úprocessr'   Úgetattr)Úr©r1   r6   r7   Ú<lambda>F   s    z*TestWorkController.start.<locals>.<lambda>)r$   r@   Ú	addFilterr+   r.   Ú
addHandlerr   r0   )r1   ÚhandlerÚloggerr4   rD   r7   r0   C   s   

zTestWorkController.startc                 C   s    | j  ¡  tj| j| |d� dS )z=Callback called when the Consumer blueprint is fully started.)Úsenderr   r   N)r   Úsetr   ÚsendÚapp)r1   r   r6   r6   r7   Úon_consumer_readyK   s   

ÿz$TestWorkController.on_consumer_readyc                 C   s   | j  ¡  dS )z±Wait for worker to be fully up and running.

        Warning:
            Worker must be started within a thread for this to work,
            or it will block forever.
        N)r   ÚwaitrD   r6   r6   r7   Úensure_startedR   s   z!TestWorkController.ensure_started)r>   r!   r?   Ú__doc__r$   r   r+   r,   r@   r0   rN   rP   Ú__classcell__r6   r6   r4   r7   r      s    
r   é   ÚsoloTg      $@c              
   k   sÞ   � t j| d� d}	z]t| f||||||dœ|¤Ž�2}	|rAddlm}
 tƒ � |
 ¡ j|d�dks2J ‚W d  ƒ n1 s<w   Y  |	V  W d  ƒ n1 sNw   Y  W tj| |	d� dS W tj| |	d� dS tj| |	d� w )	z[Start embedded worker.

    Yields:
        celery.app.worker.Worker: worker instance.
    )rJ   N)ÚconcurrencyÚpoolÚloglevelÚlogfileÚperform_ping_checkÚshutdown_timeoutrS   )Úping)ÚtimeoutÚpong)rJ   r   )	r   rL   Ú_start_worker_threadÚtasksr[   r
   ÚdelayÚgetr   )rM   rU   rV   rW   rX   rY   Úping_task_timeoutrZ   r3   r   r[   r6   r6   r7   Ústart_worker]   s2   €úùÿóñ"rc   rM   rU   rV   rW   rX   ÚWorkControllerrY   rZ   Úreturnc                 k   s&  � t | ||ƒ |rd| jv sJ ‚| jtj d¡d��}	|	jj W d  ƒ n1 s)w   Y  |d| |tƒ |||d| 	dd¡dddœ
|¤Ž}
t
j|
jdd�}| ¡  |
 ¡  td	ƒ z|
V  W d
dlm} d
|_| |¡ | ¡ rttdƒ‚d|_dS d
dlm} d
|_| |¡ | ¡ r�tdƒ‚d|_w )zaStart Celery worker in a thread.

    Yields:
        celery.worker.Worker: worker instance.
    zcelery.pingÚTEST_BROKER)ÚhostnameNÚwithout_heartbeatT)
rM   rU   rg   rV   rW   rX   Úready_callbackrh   Úwithout_mingleÚwithout_gossip)ÚtargetÚdaemonFr   )Ústatez„Worker thread failed to exit within the allocated timeout. Consider raising `shutdown_timeout` if your tasks take longer to execute.r6   )Úsetup_app_for_workerr_   Ú
connectionr%   Úenvironra   Údefault_channelÚqueue_declarer   Úpopr   ÚThreadr0   rP   r	   Úcelery.workerrn   Úshould_terminateÚjoinÚis_aliveÚRuntimeError)rM   rU   rV   rW   rX   rd   rY   rZ   r3   Úconnr   Útrn   r6   r6   r7   r^   …   sV   €
ÿ
õô
ÿ
÷
ÿr^   c           	      k   sP   � ddl m}m} |  ¡  ||dƒgƒ}| ¡  z
dV  W | ¡  dS | ¡  w )zfStart worker in separate process.

    Yields:
        celery.app.worker.Worker: worker instance.
    r   )ÚClusterÚNodeztestworker1@%hN)Úcelery.apps.multir}   r~   Úset_currentr0   Ústopwait)	rM   rU   rV   rW   rX   r3   r}   r~   Úclusterr6   r6   r7   Ú_start_worker_process½   s   €rƒ   c                 C   s8   |   ¡  |  ¡  |  ¡  dt| jƒ_| jj||d� dS )z9Setup the app to be used for starting an embedded worker.F)rW   rX   N)Úfinalizer€   Úset_defaultÚtypeÚlogÚ_setupÚsetup)rM   rW   rX   r6   r6   r7   ro   Õ   s
   ro   )&rQ   r+   r%   r   Ú
contextlibr   Útypingr   r   r   r   Úcelery.worker.consumerÚceleryr   r   Úcelery.resultr	   r
   Úcelery.utils.dispatchr   Úcelery.utils.nodenamesr   rq   ra   r   r   r   r   rd   r   rc   ÚintÚstrÚboolÚfloatr^   rƒ   ro   r6   r6   r6   r7   Ú<module>   s„    þþþ?ø'ùÿþ
ýüûúùø7ü&