o
    wvXjI8  ã                   @   s¤  d Z ddlmZmZm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mZmZ ddlmZ ddlmZ dZdZG dd„ deƒZeG dd„ de ƒƒZ!		d6dd„Z"d7dd„Z#dd„ Z$e#ddfdd„Z%dd„ Z&		d8dd„Z'dd „ Z(d!d"„ Z)d#d$„ Z*d%d&„ Z+G d'd(„ d(e ƒZ,	)			d9d+d,„Z-d-d.„ Z.d/d0„ Z/d1d2„ Z0d3d4„ Z1ee'ed5�Z2ee.ed5�Z3ee/ed5�Z4ee0ed5�Z5dS ):z,Message migration tools (Broker <-> Broker).é    )Úabsolute_importÚprint_functionÚunicode_literalsN)Úpartial)ÚcycleÚislice)ÚQueueÚ	eventloop)Úmaybe_declare)Úensure_bytes)Úapp_or_default)Úpython_2_unicode_compatibleÚstringÚstring_t)Ú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$   úS/var/www/html/myproject/venv/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_apxr   ©Úselfr$   r$   r%   Ústrtotal+   s   
zState.strtotalc                 C   s   | j rd | ¡S d | ¡S )Nz^{0.filtered}z{0.count}/{0.strtotal})ÚfilteredÚformatr(   r$   r$   r%   Ú__repr__1   s   

zState.__repr__N)
r    r!   r"   r#   Úcountr+   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¡ q9| jt |ƒf|||||	|
dœ|¤Ž dS )zRepublish message.)Úapplication_headersÚcontent_typeÚcontent_encodingÚheadersNÚexchangeÚrouting_keyÚcompression)r4   r5   r6   r3   r1   r2   )	r   ÚbodyÚdelivery_infor3   Ú
propertiesr1   r2   ÚpopÚpublish)ÚproducerÚmessager4   r5   Úremove_propsr7   Úinfor3   ÚpropsÚctypeÚencr6   Úkeyr$   r$   r%   r   7   s&   

ÿý
ýr   c                 C   s>   |j }|du r	i n|}t| || |d ¡| |d ¡d� dS )zMigrate single task message.Nr4   r5   ©r4   r5   )r8   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$   ©r7   r=   ©ÚcallbackÚtasksr$   r%   r+   [   s   
z!filter_callback.<locals>.filteredr$   )rL   rM   r+   r$   rK   r%   Úfilter_callbackY   s   rN   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©rG   c                    sh   | ˆ j ƒ}ˆ | j| j¡|_|j| jkrˆ | j|j¡|_|jj| jkr.ˆ | j| j¡|j_| ¡  d S ©N)ÚchannelrE   Únamer5   r4   Údeclare)ÚqueueÚ	new_queue©r<   rG   r$   r%   Úon_declare_queuek   s   
ÿz'migrate_tasks.<locals>.on_declare_queue)rG   rX   )r   Úprepare_queuesÚamqpÚProducerr   r   )ÚsourceÚdestÚmigrateÚapprG   ÚkwargsrX   r$   rW   r%   r   c   s   
ÿÿr   c                 C   s   t |tƒr| jj| S |S rQ   )Ú
isinstancer   rZ   rG   )r_   Úqr$   r$   r%   Ú_maybe_queuey   s   
rc   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 )	aL	  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, string_t):
                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$   )rc   )Ú.0rU   )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 )NrD   é   )ra   r   r
   Údefault_channelr4   rS   r5   Úexpand_destr   Úackr+   r   )r7   r=   ÚretÚexÚrk)	rL   Úconnr4   ÚlimitÚ	predicater<   r5   ÚstateÚ	transformr$   r%   Úon_taskÃ   s&   

ÿðzmove.<locals>.on_taskÚconsume_from)r   Úconnection_or_acquirerZ   r[   r   r   )rp   Ú
connectionr4   r5   r\   r_   rL   ro   rr   r`   rG   rs   r$   )
r_   rL   rn   r4   ro   rp   r<   r5   rq   rr   r%   r      s   >$èr   c              	   C   s:   z	| \}}W ||fS  t tfy   ||}}Y ||fS w rQ   )Ú	TypeErrorÚ
ValueError)rk   r4   r5   rl   rm   r$   r$   r%   ri   Ú   s   
þþri   c                 C   s   |d | kS )z'Return true if task id equals task_id'.Úidr$   )Útask_idr7   r=   r$   r$   r%   r   â   ó   r   c                 C   s   |d | v S )z-Return true if task id is member of set ids'.ry   r$   )Úidsr7   r=   r$   r$   r%   r   ç   r{   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©rd   rb   r$   r$   r%   Ú	<genexpr>ð   s   € "ÿz!prepare_queues.<locals>.<genexpr>)ra   r   r�   ÚlistÚdictrP   r$   r$   r%   rY   ì   s   


ÿrY   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$   )rc   r_   r‚   r(   r$   r%   re   	  s    
ÿÿz%Filterer.__init__.<locals>.<listcomp>)r_   rn   Úfilterro   ÚtimeoutÚack_messagesÚsetr   rM   rY   rG   rL   ÚforeverrX   r„   rt   r   rq   Úaccept)r)   r_   rn   rˆ   ro   r‰   rŠ   rM   rG   rL   rŒ   rX   rt   rq   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	   rn   r‰   rŒ   Úsocketr   rq   )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 )Nrg   )rq   r.   ro   r   ©r)   r7   r=   r$   r$   r%   Úupdate_state  s   ÿzFilterer.update_statec                 C   s   |  ¡  d S rQ   )rj   r•   r$   r$   r%   Úack_message#  s   zFilterer.ack_messagec                 C   s   | j jj| j| j| jd�S )N)rG   r�   )r_   rZ   ÚTaskConsumerrn   rt   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 rQ   )rˆ   r–   r—   rM   rN   Úregister_callbackrŠ   rL   r   rq   Údeclare_queues)r)   Úconsumerrˆ   r–   r—   rL   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)	rG   rS   rX   rR   Úqueue_declarerq   r'   rn   Úchannel_errors)r)   r›   rU   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.)ro   r‰   rŠ   rM   rG   rL   rŒ   rX   rt   rq   r�   )r†   r”   )r_   rn   rˆ   ro   r‰   rŠ   rM   rG   rL   rŒ   rX   rt   rq   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   )rz   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)rE   r9   rJ   ©Úmapr$   r%   Útask_id_in_map€  s   z%move_by_idmap.<locals>.task_id_in_mapro   )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 rH   )rE   rJ   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)rq   r7   r$   )ÚprintÚMOVING_PROGRESS_FMTr,   )rq   r7   r=   r`   r$   r$   r%   Úfilter_status™  s   r©   )rr   )NNNrQ   )NNNNNNNNr    )6r#   Ú
__future__r   r   r   r’   Ú	functoolsr   Ú	itertoolsr   r   Úkombur   r	   Úkombu.commonr
   Úkombu.utils.encodingr   Ú
celery.appr   Úcelery.fiver   r   r   Úcelery.utils.nodenamesr   Úcelery.utils.textr   Ú__all__r¨   Ú	Exceptionr   Úobjectr   r   r   rN   r   rc   r   ri   r   r   rY   r†   r   r   r   r   r©   r   r   Úmove_direct_by_idmapÚmove_direct_by_taskmapr$   r$   r$   r%   Ú<module>   s^   
ÿ
	

ÿ
ÿ[Z
ý