o
    ›¨Êh8  ã                   @   st  d Z ddlZddlm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 ddlmZ dd	lmZ dd
lmZ dZdZG dd„ deƒZG dd„ dƒZ		d4dd„Zd5dd„Zdd„ Zeddfdd„Zdd„ Z		d6dd„Zdd„ Zdd „ Z d!d"„ Z!d#d$„ Z"G d%d&„ d&ƒZ#	'			d7d)d*„Z$d+d,„ Z%d-d.„ Z&d/d0„ Z'd1d2„ Z(eeed3�Z)ee%ed3�Z*ee&ed3�Z+ee'ed3�Z,dS )8z,Message migration tools (Broker <-> Broker).é    N)Úpartial)ÚcycleÚislice)ÚQueueÚ	eventloop)Úmaybe_declare)Úensure_bytes)Úapp_or_default)Úworker_direct)Ústr_to_list)ÚStopFilteringÚStateÚ	republishÚmigrate_taskÚmigrate_tasksÚmoveÚ
task_id_eqÚ
task_id_inÚstart_filterÚmove_task_by_idÚmove_by_idmapÚmove_by_taskmapÚmove_directÚmove_direct_by_idzGMoving task {state.filtered}/{state.strtotal}: {body[task]}[{body[id]}]c                   @   s   e Zd ZdZdS )r   z*Semi-predicate used to signal filter stop.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__© r   r   úH/var/www/html/env/lib/python3.10/site-packages/celery/contrib/migrate.pyr      s    r   c                   @   s0   e Zd ZdZdZdZdZedd„ ƒZdd„ Z	dS )r   zMigration progress state.r   c                 C   s   | j sdS t| j ƒS )Nú?)Ú	total_apxÚstr©Úselfr   r   r   Ústrtotal&   s   
zState.strtotalc                 C   s$   | j r	d| j › �S | j› d| j› �S )Nú^ú/)ÚfilteredÚcountr%   r#   r   r   r   Ú__repr__,   s   zState.__repr__N)
r   r   r   r   r)   r(   r!   Úpropertyr%   r*   r   r   r   r   r      s    
r   c              
   C   sÎ   |sg d¢}t |jƒ}|j|j|j}}}|du r|d n|}|du r(|d n|}|j|j}	}
| dd¡}| dd¡}|durEt|ƒnd}|D ]}| |d¡ qI| j	t |ƒf|||||	|
|dœ|¤Ž dS )zRepublish message.)Úapplication_headersÚcontent_typeÚcontent_encodingÚheadersNÚexchangeÚrouting_keyÚcompressionÚ
expiration)r0   r1   r2   r/   r-   r.   r3   )
r   ÚbodyÚdelivery_infor/   Ú
propertiesr-   r.   ÚpopÚfloatÚpublish)ÚproducerÚmessager0   r1   Úremove_propsr4   Úinfor/   ÚpropsÚctypeÚencr2   r3   Úkeyr   r   r   r   2   s*   

ÿý
ür   c                 C   s>   |j }|du r	i n|}t| || |d ¡| |d ¡d� dS )zMigrate single task message.Nr0   r1   ©r0   r1   )r5   r   Úget)r:   Úbody_r;   Úqueuesr=   r   r   r   r   P   s   
þr   c                    s   ‡ ‡fdd„}|S )Nc                    s   ˆr
| d ˆvr
d S ˆ | |ƒS ©NÚtaskr   ©r4   r;   ©ÚcallbackÚtasksr   r   r(   [   s   
z!filter_callback.<locals>.filteredr   )rJ   rK   r(   r   rI   r   Úfilter_callbackY   s   rL   c                    sV   t |ƒ}tˆƒ‰|jj|dd�‰ t|ˆ ˆd�}‡ ‡fdd„}t|| |fˆ|dœ|¤ŽS )z)Migrate tasks from one broker to another.F)Úauto_declare©rE   c                    sh   | ˆ j ƒ}ˆ | j| j¡|_|j| jkrˆ | j|j¡|_|jj| jkr.ˆ | j| j¡|j_| ¡  d S ©N)ÚchannelrC   Únamer1   r0   Údeclare)ÚqueueÚ	new_queue©r:   rE   r   r   Úon_declare_queuek   s   
ÿz'migrate_tasks.<locals>.on_declare_queue)rE   rV   )r	   Úprepare_queuesÚamqpÚProducerr   r   )ÚsourceÚdestÚmigrateÚapprE   ÚkwargsrV   r   rU   r   r   c   s   
ÿÿr   c                 C   s   t |tƒr| jj| S |S rO   )Ú
isinstancer"   rX   rE   )r]   Úqr   r   r   Ú_maybe_queuey   s   
ra   c	              
      sš   t ˆ ƒ‰ ‡ fdd„|pg D ƒpd}
ˆ j|dd��+‰ˆ j ˆ¡‰tƒ ‰‡‡‡‡‡‡‡‡‡	f	dd„}tˆ ˆ|fd|
i|	¤ŽW  d  ƒ S 1 sFw   Y  dS )	aG	  Find tasks by filtering them and move the tasks to a new queue.

    Arguments:
        predicate (Callable): Filter function used to decide the messages
            to move.  Must accept the standard signature of ``(body, message)``
            used by Kombu consumer callbacks.  If the predicate wants the
            message to be moved it must return either:

                1) a tuple of ``(exchange, routing_key)``, or

                2) a :class:`~kombu.entity.Queue` instance, or

                3) any other true value means the specified
                    ``exchange`` and ``routing_key`` arguments will be used.
        connection (kombu.Connection): Custom connection to use.
        source: List[Union[str, kombu.Queue]]: Optional list of source
            queues to use instead of the default (queues
            in :setting:`task_queues`).  This list can also contain
            :class:`~kombu.entity.Queue` instances.
        exchange (str, kombu.Exchange): Default destination exchange.
        routing_key (str): Default destination routing key.
        limit (int): Limit number of messages to filter.
        callback (Callable): Callback called after message moved,
            with signature ``(state, body, message)``.
        transform (Callable): Optional function to transform the return
            value (destination) of the filter function.

    Also supports the same keyword arguments as :func:`start_filter`.

    To demonstrate, the :func:`move_task_by_id` operation can be implemented
    like this:

    .. code-block:: python

        def is_wanted_task(body, message):
            if body['id'] == wanted_id:
                return Queue('foo', exchange=Exchange('foo'),
                             routing_key='foo')

        move(is_wanted_task)

    or with a transform:

    .. code-block:: python

        def transform(value):
            if isinstance(value, str):
                return Queue(value, Exchange(value), value)
            return value

        move(is_wanted_task, transform=transform)

    Note:
        The predicate may also return a tuple of ``(exchange, routing_key)``
        to specify the destination to where the task should be moved,
        or a :class:`~kombu.entity.Queue` instance.
        Any other true value means that the task will be moved to the
        default exchange/routing_key.
    c                    s   g | ]}t ˆ |ƒ‘qS r   )ra   )Ú.0rS   )r]   r   r   Ú
<listcomp>¾   s    zmove.<locals>.<listcomp>NF)Úpoolc                    s¨   ˆ| |ƒ}|rNˆrˆ|ƒ}t |tƒr!t|ˆjƒ |jj|j}}nt|ˆˆƒ\}}tˆ|||d� | 	¡  ˆ j
d7  _
ˆ rDˆ ˆ| |ƒ ˆrPˆj
ˆkrRtƒ ‚d S d S d S )NrB   é   )r_   r   r   Údefault_channelr0   rQ   r1   Úexpand_destr   Úackr(   r   )r4   r;   ÚretÚexÚrk)	rJ   Úconnr0   ÚlimitÚ	predicater:   r1   ÚstateÚ	transformr   r   Úon_taskÃ   s&   

ÿðzmove.<locals>.on_taskÚconsume_from)r	   Úconnection_or_acquirerX   rY   r   r   )rn   Ú
connectionr0   r1   rZ   r]   rJ   rm   rp   r^   rE   rq   r   )
r]   rJ   rl   r0   rm   rn   r:   r1   ro   rp   r   r      s   >$èr   c              	   C   s:   z	| \}}W ||fS  t tfy   ||}}Y ||fS w rO   )Ú	TypeErrorÚ
ValueError)ri   r0   r1   rj   rk   r   r   r   rg   Ú   s   
þþrg   c                 C   s   |d | kS )z'Return true if task id equals task_id'.Úidr   )Útask_idr4   r;   r   r   r   r   â   ó   r   c                 C   s   |d | v S )z-Return true if task id is member of set ids'.rw   r   )Úidsr4   r;   r   r   r   r   ç   ry   r   c                 C   s@   t | tƒr
|  d¡} t | tƒrtdd„ | D ƒƒ} | d u ri } | S )Nú,c                 s   s*   � | ]}t tt| d ¡ƒddƒƒV  qdS )ú:Né   )Útupler   r   Úsplit©rb   r`   r   r   r   Ú	<genexpr>ð   s   € "ÿz!prepare_queues.<locals>.<genexpr>)r_   r"   r   ÚlistÚdictrN   r   r   r   rW   ì   s   


ÿrW   c                   @   sN   e Zd Z				ddd„Zdd„ Zdd	„ Zd
d„ Zdd„ Zdd„ Zdd„ Z	dS )ÚFiltererNç      ð?Fc                    s†   |ˆ _ |ˆ _|ˆ _|ˆ _|ˆ _|ˆ _tt|ƒpg ƒˆ _t	|ƒˆ _
|	ˆ _|
ˆ _|ˆ _‡ fdd„|p4tˆ j
ƒD ƒˆ _|p<tƒ ˆ _|ˆ _d S )Nc                    s   g | ]}t ˆ j|ƒ‘qS r   )ra   r]   r€   r#   r   r   rc   	  s    
ÿÿz%Filterer.__init__.<locals>.<listcomp>)r]   rl   Úfilterrm   ÚtimeoutÚack_messagesÚsetr   rK   rW   rE   rJ   ÚforeverrV   r‚   rr   r   ro   Úaccept)r$   r]   rl   r†   rm   r‡   rˆ   rK   rE   rJ   rŠ   rV   rr   ro   r‹   r^   r   r#   r   Ú__init__ù   s    

þ
zFilterer.__init__c              	   C   s    |   |  ¡ ¡�> zt| j| j| jd�D ]}qW n tjy!   Y n ty)   Y nw W d   ƒ | jS W d   ƒ | jS W d   ƒ | jS 1 sHw   Y  | jS )N)r‡   Úignore_timeouts)	Úprepare_consumerÚcreate_consumerr   rl   r‡   rŠ   Úsocketr   ro   )r$   Ú_r   r   r   Ústart  s0   
þýÿú
þ
ý
ù
ÿ
÷
ö
zFilterer.startc                 C   s2   | j  jd7  _| jr| j j| jkrtƒ ‚d S d S )Nre   )ro   r)   rm   r   ©r$   r4   r;   r   r   r   Úupdate_state  s   ÿzFilterer.update_statec                 C   s   |  ¡  d S rO   )rh   r“   r   r   r   Úack_message#  s   zFilterer.ack_messagec                 C   s   | j jj| j| j| jd�S )N)rE   r‹   )r]   rX   ÚTaskConsumerrl   rr   r‹   r#   r   r   r   r�   &  s
   ýzFilterer.create_consumerc                 C   s¤   | j }| j}| j}| jrt|| jƒ}t|| jƒ}t|| jƒ}| |¡ | |¡ | jr1| | j¡ | jd urKt| j| j	ƒ}| jrFt|| jƒ}| |¡ |  
|¡ |S rO   )r†   r”   r•   rK   rL   Úregister_callbackrˆ   rJ   r   ro   Údeclare_queues)r$   Úconsumerr†   r”   r•   rJ   r   r   r   rŽ   -  s$   




zFilterer.prepare_consumerc              	   C   s~   |j D ]9}| j r|j| j vrq| jd ur|  |¡ z||jƒjdd�\}}}|r0| j j|7  _W q | jjy<   Y qw d S )NT)Úpassive)	rE   rQ   rV   rP   Úqueue_declarero   r!   rl   Úchannel_errors)r$   r™   rS   r‘   Úmcountr   r   r   r˜   A  s$   


ÿÿ€ÿözFilterer.declare_queues©Nr…   FNNNFNNNN)
r   r   r   rŒ   r’   r”   r•   r�   rŽ   r˜   r   r   r   r   r„   ÷   s    
ür„   r…   Fc                 K   s0   t | ||f|||||||	|
|||dœ|¤Ž ¡ S )zFilter tasks.)rm   r‡   rˆ   rK   rE   rJ   rŠ   rV   rr   ro   r‹   )r„   r’   )r]   rl   r†   rm   r‡   rˆ   rK   rE   rJ   rŠ   rV   rr   ro   r‹   r^   r   r   r   r   Q  s&   ÿôóór   c                 K   s   t | |ifi |¤ŽS )a‹  Find a task by id and move it to another queue.

    Arguments:
        task_id (str): Id of task to find and move.
        dest: (str, kombu.Queue): Destination queue.
        transform (Callable): Optional function to transform the return
            value (destination) of the filter function.
        **kwargs (Any): Also supports the same keyword
            arguments as :func:`move`.
    )r   )rx   r[   r^   r   r   r   r   f  s   r   c                    s$   ‡ fdd„}t |fdtˆ ƒi|¤ŽS )a“  Move tasks by matching from a ``task_id: queue`` mapping.

    Where ``queue`` is a queue to move the task to.

    Example:
        >>> move_by_idmap({
        ...     '5bee6e82-f4ac-468e-bd3d-13e8600250bc': Queue('name'),
        ...     'ada8652d-aef3-466b-abd2-becdaf1b82b3': Queue('name'),
        ...     '3a2b140d-7db1-41ba-ac90-c36a0ef4ab1f': Queue('name')},
        ...   queues=['hipri'])
    c                    s   ˆ   |jd ¡S )NÚcorrelation_id)rC   r6   rH   ©Úmapr   r   Útask_id_in_map€  s   z%move_by_idmap.<locals>.task_id_in_maprm   )r   Úlen)r¡   r^   r¢   r   r    r   r   t  s   r   c                    s   ‡ fdd„}t |fi |¤ŽS )a  Move tasks by matching from a ``task_name: queue`` mapping.

    ``queue`` is the queue to move the task to.

    Example:
        >>> move_by_taskmap({
        ...     'tasks.add': Queue('name'),
        ...     'tasks.mul': Queue('name'),
        ... })
    c                    s   ˆ   | d ¡S rF   )rC   rH   r    r   r   Útask_name_in_map“  s   z)move_by_taskmap.<locals>.task_name_in_map)r   )r¡   r^   r¤   r   r    r   r   ˆ  s   r   c                 K   s   t tjd| |dœ|¤Žƒ d S )N)ro   r4   r   )ÚprintÚMOVING_PROGRESS_FMTÚformat)ro   r4   r;   r^   r   r   r   Úfilter_status™  s   r¨   )rp   )NNNrO   )NNNNNNNNrž   )-r   r�   Ú	functoolsr   Ú	itertoolsr   r   Úkombur   r   Úkombu.commonr   Úkombu.utils.encodingr   Ú
celery.appr	   Úcelery.utils.nodenamesr
   Úcelery.utils.textr   Ú__all__r¦   Ú	Exceptionr   r   r   r   rL   r   ra   r   rg   r   r   rW   r„   r   r   r   r   r¨   r   r   Úmove_direct_by_idmapÚmove_direct_by_taskmapr   r   r   r   Ú<module>   sX    
ÿ
	

ÿ
ÿ[Z
ý