o
    wvXj#  ã                   @   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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fdd„ƒZeddededfdd„ƒZeddedfdd„ƒZdd„ ZdS )z'Embedded workers for integration tests.é    )Úabsolute_importÚunicode_literalsN)Úcontextmanager)Ú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                       s0   e Zd ZdZ‡ fdd„Zdd„ Zdd„ Z‡  ZS )ÚTestWorkControllerz3Worker that can synchronize on being fully started.c                    s$   t  ¡ | _tt| ƒj|i |¤Ž d S )N)Ú	threadingÚEventÚ_on_startedÚsuperr   Ú__init__)ÚselfÚargsÚkwargs©Ú	__class__© úZ/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/contrib/testing/worker.pyr       s   
zTestWorkController.__init__c                 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)r   r   r   r   r   Úon_consumer_ready%   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   Úwait)r   r   r   r   Úensure_started,   s   z!TestWorkController.ensure_started)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r#   r%   Ú__classcell__r   r   r   r   r      s
    r   é   ÚsoloTg      $@c           
   	   k   s°   � t j| d� t| f|||||dœ|¤Ž�2}|r=ddlm}	 tƒ � |	 ¡ j|d�dks.J ‚W d  ƒ n1 s8w   Y  |V  W d  ƒ n1 sJw   Y  tj| |d� dS )	z[Start embedded worker.

    Yields:
        celery.app.worker.Worker: worker instance.
    )r   )ÚconcurrencyÚpoolÚloglevelÚlogfileÚperform_ping_checkr+   )Úping)ÚtimeoutÚpongN)r   r   )	r   r!   Ú_start_worker_threadÚtasksr2   r   ÚdelayÚgetr   )
r"   r-   r.   r/   r0   r1   Úping_task_timeoutr   r   r2   r   r   r   Ústart_worker7   s(   €ûúÿôr:   c                 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œ
|¤Ž}	t	j
|	jd�}
|
 ¡  |	 ¡  tdƒ |	V  d	d
lm} d	|_|
 d¡ d|_dS )zaStart Celery worker in a thread.

    Yields:
        celery.worker.Worker: worker instance.
    zcelery.pingÚTEST_BROKER)ÚhostnameNT)
r"   r-   r<   r.   r/   r0   Úready_callbackÚwithout_heartbeatÚwithout_mingleÚwithout_gossip)ÚtargetFr   )Ústateé
   r   )Úsetup_app_for_workerr6   Ú
connectionÚosÚenvironr8   Údefault_channelÚqueue_declarer	   r   ÚThreadÚstartr%   r   Úcelery.workerrB   Úshould_terminateÚjoin)r"   r-   r.   r/   r0   ÚWorkControllerr1   r   Úconnr   ÚtrB   r   r   r   r5   Y   s<   €
ÿõô

r5   c           	      k   sB   � ddl m}m} |  ¡  ||dƒgƒ}| ¡  dV  | ¡  dS )zfStart worker in separate process.

    Yields:
        celery.app.worker.Worker: worker instance.
    r   )ÚClusterÚNodeztestworker1@%hN)Úcelery.apps.multirR   rS   Úset_currentrK   Ústopwait)	r"   r-   r.   r/   r0   r   rR   rS   Úclusterr   r   r   Ú_start_worker_processŠ   s   €rX   c                 C   s8   |   ¡  |  ¡  |  ¡  dt| jƒ_| jj||d� dS )z9Setup the app to be used for starting an embedded worker.F)r/   r0   N)ÚfinalizerU   Úset_defaultÚtypeÚlogÚ_setupÚsetup)r"   r/   r0   r   r   r   rD       s
   rD   )r)   Ú
__future__r   r   rF   r   Ú
contextlibr   Úceleryr   Úcelery.resultr   r   Úcelery.utils.dispatchr   Úcelery.utils.nodenamesr	   rG   r8   r
   r   r   r   rO   r   r:   r5   rX   rD   r   r   r   r   Ú<module>   s\    þþþú!ú0ü