o
    wvXj¨A  ã                   @   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Zdd„ Zdd„ ZG dd„ deƒZ G dd„ deƒZ!dS )z„Worker Remote Control Client.

Client for worker remote control commands.
Server implementation is in :mod:`celery.worker.control`.
é    )Úabsolute_importÚunicode_literalsN)ÚTERM_SIGNAME©Úmatch)ÚMailbox)Úregister_after_fork)Úlazy)Úcached_property)ÚDuplicateNodenameWarning)Úitems)Ú
get_logger)Ú	pluralize)ÚInspectÚControlÚflatten_replyzˆReceived multiple replies from node {0}: {1}.
Please make sure you give each node a unique nodename using
the celery worker `-n` option.c              
      sf   i t ƒ ‰‰ | D ]}‡ ‡fdd„|D ƒ ˆ |¡ qˆ r1t tt ttˆ ƒdƒd 	t
ˆ ƒ¡¡ƒ¡ ˆS )zñFlatten node replies.

    Convert from a list of replies in this format::

        [{'a@example.com': reply},
         {'b@example.com': reply}]

    into this format::

        {'a@example.com': reply,
         'b@example.com': reply}
    c                    s   g | ]}|ˆv rˆ   |¡‘qS © )Úadd)Ú.0Úname©ÚdupesÚnodesr   úO/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/app/control.pyÚ
<listcomp>1   s    z!flatten_reply.<locals>.<listcomp>r   z, )ÚsetÚupdateÚwarningsÚwarnr   Ú	W_DUPNODEÚformatr   ÚlenÚjoinÚsorted)ÚreplyÚitemr   r   r   r   "   s   ÿÿr   c              
   C   sF   z|   ¡  W d S  ty" } ztjd|dd� W Y d }~d S d }~ww )Nzafter fork raised exception: %ré   )Úexc_info)Ú_after_forkÚ	ExceptionÚloggerÚinfo)ÚcontrolÚexcr   r   r   Ú_after_fork_cleanup_control<   s   €ÿr.   c                   @   sÈ   e Zd ZdZdZ			d.dd„Zdd„ Zdd	„ Zd
d„ Zdd„ Z	d/dd„Z
d/dd„Zd/dd„Zdd„ Zdd„ Zdd„ ZeZd/dd„Zdd„ Zdd„ Zd0d!d"„Zd/d#d$„Zd%d&„ Zd1d(d)„Zd2d,d-„ZdS )3r   zAPI for app.control.inspect.Nç      ð?c	           	      C   s:   |p| j | _ || _|| _|| _|| _|| _|| _|| _d S ©N)ÚappÚdestinationÚtimeoutÚcallbackÚ
connectionÚlimitÚpatternÚmatcher)	Úselfr2   r3   r4   r5   r1   r6   r7   r8   r   r   r   Ú__init__H   s   
zInspect.__init__c                    s`   |r.t |ƒ}| jrt| jttfƒs| | j¡S | jr,| j‰| j‰ ‡ ‡fdd„t|ƒD ƒS |S d S )Nc                    s"   i | ]\}}t |ˆˆ ƒr||“qS r   r   )r   Únoder$   ©r8   r7   r   r   Ú
<dictcomp>]   s    
ÿz$Inspect._prepare.<locals>.<dictcomp>)	r   r2   Ú
isinstanceÚlistÚtupleÚgetr7   r8   r   )r9   r$   Úby_noder   r<   r   Ú_prepareT   s   ÿözInspect._preparec                 K   s6   |   | jjj||| j| j| j| j| jd| j	| j
d�
¡S )NT)	Ú	argumentsr2   r4   r5   r6   r3   r$   r7   r8   )rC   r1   r,   Ú	broadcastr2   r4   r5   r6   r3   r7   r8   )r9   ÚcommandÚkwargsr   r   r   Ú_requesta   s   øzInspect._requestc                 C   ó
   |   d¡S )NÚreport©rH   ©r9   r   r   r   rJ   m   ó   
zInspect.reportc                 C   rI   )NÚclockrK   rL   r   r   r   rN   p   rM   zInspect.clockc                 C   rI   )NÚactiverK   ©r9   Úsafer   r   r   rO   s   s   
zInspect.activec                 C   rI   )NÚ	scheduledrK   rP   r   r   r   rR   y   rM   zInspect.scheduledc                 C   rI   )NÚreservedrK   rP   r   r   r   rS   |   rM   zInspect.reservedc                 C   rI   )NÚstatsrK   rL   r   r   r   rT      rM   zInspect.statsc                 C   rI   )NÚrevokedrK   rL   r   r   r   rU   ‚   rM   zInspect.revokedc                 G   ó   | j d|d�S )NÚ
registered)ÚtaskinfoitemsrK   )r9   rX   r   r   r   rW   …   ó   zInspect.registeredc                 C   rI   )NÚpingrK   )r9   r2   r   r   r   rZ   ‰   rM   zInspect.pingc                 C   rI   )NÚactive_queuesrK   rL   r   r   r   r[   Œ   rM   zInspect.active_queuesc                 G   s4   t |ƒdkrt|d ttfƒr|d }| jd|d�S )Nr&   r   Ú
query_task)Úids)r!   r>   r?   r@   rH   )r9   r]   r   r   r   r\   �   s   zInspect.query_taskFc                 C   rV   )NÚconf)Úwith_defaultsrK   )r9   r_   r   r   r   r^   –   rY   zInspect.confc                 C   s   | j d||d�S )NÚhello)Ú	from_noderU   rK   )r9   ra   rU   r   r   r   r`   ™   s   zInspect.helloc                 C   rI   )NÚ	memsamplerK   rL   r   r   r   rb   œ   rM   zInspect.memsampleé
   c                 C   rV   )NÚmemdump)ÚsamplesrK   )r9   re   r   r   r   rd   Ÿ   rY   zInspect.memdumpÚRequestéÈ   c                 C   s   | j d|||d�S )NÚobjgraph)ÚnumÚ	max_depthÚtyperK   )r9   rk   Únrj   r   r   r   rh   ¢   s   zInspect.objgraph)Nr/   NNNNNNr0   )F)rc   )rf   rg   rc   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r1   r:   rC   rH   rJ   rN   rO   rR   rS   rT   rU   rW   Úregistered_tasksrZ   r[   r\   r^   r`   rb   rd   rh   r   r   r   r   r   C   s4    
þ






r   c                   @   s  e Zd ZdZeZd1dd„Zdd„ Zedd„ ƒZd1d	d
„Z	e	Z
d2dd„Zddefdd„Zdefdd„Zd3dd„Zd1dd„Z		d4dd„Zd1dd„Z		d5dd„Zd1dd„Zd1d d!„Zd6d#d$„Zd6d%d&„Zd1d'd(„Zd1d)d*„Z		d7d+d,„Zd1d-d.„Z			d8d/d0„ZdS )9r   zWorker remote control client.Nc              
      sR   |ˆ _ ˆ j|jjddgt‡ fdd„ƒ|jj|jj|jj|jjd�ˆ _tˆ t	ƒ d S )NÚfanoutÚjsonc                      s
   ˆ j jjS r0   )r1   ÚamqpÚproducer_poolr   rL   r   r   Ú<lambda>±   s   
 z"Control.__init__.<locals>.<lambda>)rk   Úacceptru   Ú	queue_ttlÚreply_queue_ttlÚqueue_expiresÚreply_queue_expires)
r1   r   r^   Úcontrol_exchanger	   Úcontrol_queue_ttlÚcontrol_queue_expiresÚmailboxr   r.   )r9   r1   r   rL   r   r:   «   s   ø
zControl.__init__c                 C   s
   | j `d S r0   )r   ru   rL   r   r   r   r(   ¹   rM   zControl._after_forkc                 C   s   | j jtdd�S )Nzcontrol.inspect)Úreverse)r1   Úsubclass_with_selfr   rL   r   r   r   Úinspect¼   s   zControl.inspectc                 C   sB   | j  |¡�}| j j |¡ ¡ W  d  ƒ S 1 sw   Y  dS )a²  Discard all waiting tasks.

        This will ignore all tasks waiting for execution, and they will
        be deleted from the messaging server.

        Arguments:
            connection (kombu.Connection): Optional specific connection
                instance to use.  If not provided a connection will
                be acquired from the connection pool.

        Returns:
            int: the number of tasks discarded.
        N)r1   Úconnection_or_acquirert   ÚTaskConsumerÚpurge)r9   r5   Úconnr   r   r   r…   À   s   $ÿzControl.purgec                 C   s   | j d|d |||dœd� d S )NÚelection)ÚidÚtopicÚaction)r5   r2   rD   ©rE   )r9   rˆ   r‰   rŠ   r5   r   r   r   r‡   Ò   s
   ÿ
þzControl.electionFc                 K   s   | j d||||dœdœ|¤ŽS )a[  Tell all (or specific) workers to revoke a task by id (or list of ids).

        If a task is revoked, the workers will ignore the task and
        not execute it after all.

        Arguments:
            task_id (Union(str, list)): Id of the task to revoke
                (or list of ids).
            terminate (bool): Also terminate the process currently working
                on the task (if any).
            signal (str): Name of signal to send to process if terminate.
                Default is TERM.

        See Also:
            :meth:`broadcast` for supported keyword arguments.
        Úrevoke)Útask_idÚ	terminateÚsignal©r2   rD   N©rŒ   r‹   )r9   r�   r2   rŽ   r�   rG   r   r   r   rŒ   Ú   s   ýüzControl.revokec                 K   s   | j |f|d|dœ|¤ŽS )zÍTell all (or specific) workers to terminate a task by id (or list of ids).

        See Also:
            This is just a shortcut to :meth:`revoke` with the terminate
            argument enabled.
        T)r2   rŽ   r�   r‘   )r9   r�   r2   r�   rG   r   r   r   rŽ   ò   s   ÿþþzControl.terminater/   c                 K   s   | j 	ddi ||dœ|¤ŽS )zÒPing all (or specific) workers.

        Returns:
            List[Dict]: List of ``{'hostname': reply}`` dictionaries.

        See Also:
            :meth:`broadcast` for supported keyword arguments.
        rZ   T)r$   rD   r2   r3   N)rZ   r‹   )r9   r2   r3   rG   r   r   r   rZ   þ   s   	ÿþþzControl.pingc                 K   s   | j 	d|||dœdœ|¤ŽS )aÌ  Tell workers to set a new rate limit for task by type.

        Arguments:
            task_name (str): Name of task to change rate limit for.
            rate_limit (int, str): The rate limit as tasks per second,
                or a rate limit string (`'100/m'`, etc.
                see :attr:`celery.task.base.Task.rate_limit` for
                more information).

        See Also:
            :meth:`broadcast` for supported keyword arguments.
        Ú
rate_limit)Ú	task_namer’   r�   N)r’   r‹   )r9   r“   r’   r2   rG   r   r   r   r’     s   ÿþýùzControl.rate_limitÚdirectc              	   K   s2   | j 	d|t||||dœfi |pi ¤Ždœ|¤ŽS )aŽ  Tell all (or specific) workers to start consuming from a new queue.

        Only the queue name is required as if only the queue is specified
        then the exchange/routing key will be set to the same name (
        like automatic queues do).

        Note:
            This command does not respect the default queue/exchange
            options in the configuration.

        Arguments:
            queue (str): Name of queue to start consuming from.
            exchange (str): Optional name of exchange.
            exchange_type (str): Type of exchange (defaults to 'direct')
                command to, when empty broadcast to all workers.
            routing_key (str): Optional routing key.
            options (Dict): Additional options as supported
                by :meth:`kombu.entity.Queue.from_dict`.

        See Also:
            :meth:`broadcast` for supported keyword arguments.
        Úadd_consumer)ÚqueueÚexchangeÚexchange_typeÚrouting_keyr�   N)r•   )rE   Údict)r9   r–   r—   r˜   r™   Úoptionsr2   rG   r   r   r   r•   !  s   ÿüûý	÷zControl.add_consumerc                 K   s   | j 	d|d|idœ|¤ŽS )zšTell all (or specific) workers to stop consuming from ``queue``.

        See Also:
            Supports the same arguments as :meth:`broadcast`.
        Úcancel_consumerr–   r�   N)rœ   r‹   )r9   r–   r2   rG   r   r   r   rœ   F  s   ÿþþzControl.cancel_consumerc                 K   s    | j 	d|||dœ|dœ|¤ŽS )aS  Tell workers to set time limits for a task by type.

        Arguments:
            task_name (str): Name of task to change time limits for.
            soft (float): New soft time limit (in seconds).
            hard (float): New hard time limit (in seconds).
            **kwargs (Any): arguments passed on to :meth:`broadcast`.
        Ú
time_limit)r“   ÚhardÚsoft©rD   r2   N)r�   r‹   )r9   r“   rŸ   rž   r2   rG   r   r   r   r�   P  s   
ÿýùøzControl.time_limitc                 K   ó   | j 	di |dœ|¤ŽS )zŠTell all (or specific) workers to enable events.

        See Also:
            Supports the same arguments as :meth:`broadcast`.
        Úenable_eventsr    N)r¢   r‹   ©r9   r2   rG   r   r   r   r¢   d  ó   ÿÿÿzControl.enable_eventsc                 K   r¡   )z‹Tell all (or specific) workers to disable events.

        See Also:
            Supports the same arguments as :meth:`broadcast`.
        Údisable_eventsr    N)r¥   r‹   r£   r   r   r   r¥   m  r¤   zControl.disable_eventsr&   c                 K   ó   | j 	dd|i|dœ|¤ŽS )z“Tell all (or specific) workers to grow the pool by ``n``.

        See Also:
            Supports the same arguments as :meth:`broadcast`.
        Ú	pool_growrl   r    N)r§   r‹   ©r9   rl   r2   rG   r   r   r   r§   v  s   ÿÿÿzControl.pool_growc                 K   r¦   )z•Tell all (or specific) workers to shrink the pool by ``n``.

        See Also:
            Supports the same arguments as :meth:`broadcast`.
        Úpool_shrinkrl   r    N)r©   r‹   r¨   r   r   r   r©     s   ÿþþzControl.pool_shrinkc                 K   s   | j 	d||dœ|dœ|¤ŽS )z}Change worker(s) autoscale setting.

        See Also:
            Supports the same arguments as :meth:`broadcast`.
        Ú	autoscale)ÚmaxÚminr    N)rª   r‹   )r9   r«   r¬   r2   rG   r   r   r   rª   ‰  s   ÿþþzControl.autoscalec                 K   r¡   )zlShutdown worker(s).

        See Also:
            Supports the same arguments as :meth:`broadcast`
        Úshutdownr    N)r­   r‹   r£   r   r   r   r­   “  r¤   zControl.shutdownc                 K   s    | j 	d|||dœ|dœ|¤ŽS )aÛ  Restart the execution pools of all or specific workers.

        Keyword Arguments:
            modules (Sequence[str]): List of modules to reload.
            reload (bool): Flag to enable module reloading.  Default is False.
            reloader (Any): Function to reload a module.
            destination (Sequence[str]): List of worker names to send this
                command to.

        See Also:
            Supports the same arguments as :meth:`broadcast`
        Úpool_restart)ÚmodulesÚreloadÚreloaderr    N)r®   r‹   )r9   r¯   r°   r±   r2   rG   r   r   r   r®   œ  s   ÿýùùzControl.pool_restartc                 K   r¡   )zˆTell worker(s) to send a heartbeat immediately.

        See Also:
            Supports the same arguments as :meth:`broadcast`
        Ú	heartbeatr    N)r²   r‹   r£   r   r   r   r²   ³  r¤   zControl.heartbeatc                 K   sž   | j  |¡�?}t|pi fi |¤Ž}|
r.|r.|  |¡j||||||||	|
|d�
W  d  ƒ S |  |¡j||||||||	d�W  d  ƒ S 1 sHw   Y  dS )a  Broadcast a control command to the celery workers.

        Arguments:
            command (str): Name of command to send.
            arguments (Dict): Keyword arguments for the command.
            destination (List): If set, a list of the hosts to send the
                command to, when empty broadcast to all workers.
            connection (kombu.Connection): Custom broker connection to use,
                if not set, a connection will be acquired from the pool.
            reply (bool): Wait for and return the reply.
            timeout (float): Timeout in seconds to wait for the reply.
            limit (int): Limit number of replies.
            callback (Callable): Callback called immediately for
                each reply received.
            pattern (str): Custom pattern string to match
            matcher (Callable): Custom matcher to run the pattern to match
        )Úchannelr7   r8   N)r³   )r1   rƒ   rš   r   Ú
_broadcast)r9   rF   rD   r2   r5   r$   r3   r6   r4   r³   r7   r8   Úextra_kwargsr†   r   r   r   rE   ¼  s   

ýû

þ$õzControl.broadcastr0   )NN)Nr/   )Nr”   NNN)NNN)r&   N)NFNN)
NNNFr/   NNNNN)rm   rn   ro   rp   r   r:   r(   r
   r‚   r…   Údiscard_allr‡   r   rŒ   rŽ   rZ   r’   r•   rœ   r�   r¢   r¥   r§   r©   rª   r­   r®   r²   rE   r   r   r   r   r   ¦   sL    




ÿ
ÿ


þ
%

ÿ

	
	
	



	
ÿ
	þr   )"rp   Ú
__future__r   r   r   Úbilliard.commonr   Úkombu.matcherr   Úkombu.pidboxr   Úkombu.utils.compatr   Úkombu.utils.functionalr	   Úkombu.utils.objectsr
   Úcelery.exceptionsr   Úcelery.fiver   Úcelery.utils.logr   Úcelery.utils.textr   Ú__all__rm   r*   r   r   r.   Úobjectr   r   r   r   r   r   Ú<module>   s(   c