o
    ›¨ÊhÄZ  ã                   @   s
  d 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mZmZmZ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dlmZ ddlmZ ddl m!Z! ddl"m#Z$ dZ%dZ&dZ'eddƒZ(ddd„Z)G dd„ de*ƒZ+G dd„ dƒZ,dS )z/Sending/Receiving Messages (Kombu integration).é    N)Ú
namedtuple)ÚMapping)Ú	timedelta)ÚWeakValueDictionary)Ú
ConnectionÚConsumerÚExchangeÚProducerÚQueueÚpools)Ú	Broadcast)Ú
maybe_list)Úcached_property)Úsignals)Úanon_nodename)Úsaferepr)Úindent)Úmaybe_make_awareé   )Úroutes)ÚAMQPÚQueuesÚtask_messagei   €zS
.> {0.name:<16} exchange={0.exchange.name}({0.exchange.type}) key={0.routing_key}
r   ©ÚheadersÚ
propertiesÚbodyÚ
sent_eventúutf-8c                    s   ‡ fdd„|   ¡ D ƒS )Nc                    s*   i | ]\}}t |tƒr| ˆ ¡n||“qS © )Ú
isinstanceÚbytesÚdecode)Ú.0ÚkÚv©Úencodingr   úA/var/www/html/env/lib/python3.10/site-packages/celery/app/amqp.pyÚ
<dictcomp>%   s    ÿzutf8dict.<locals>.<dictcomp>)Úitems)Údr'   r   r&   r(   Úutf8dict$   s   
ÿr,   c                       s¢   e Zd ZdZdZ			d!‡ fdd„	Z‡ fdd„Z‡ fdd	„Zd
d„ Zdd„ Z	dd„ Z
dd„ Zdd„ Zd"dd„Zdd„ Zdd„ Zdd„ Zdd„ Zedd „ ƒZ‡  ZS )#r   u�  Queue nameâ‡’ declaration mapping.

    Arguments:
        queues (Iterable): Initial list/tuple or dict of queues.
        create_missing (bool): By default any unknown queues will be
            added automatically, but if this flag is disabled the occurrence
            of unknown queues in `wanted` will raise :exc:`KeyError`.
        max_priority (int): Default x-max-priority for queues with none set.
    NTc           	         s    t ƒ  ¡  tƒ | _|| _|| _|| _|d u rtn|| _|| _	|d ur.t
|tƒs.dd„ |D ƒ}|p1i }| ¡ D ]\}}t
|tƒrD|  |¡n| j|fi |¤Ž q6d S )Nc                 S   s   i | ]}|j |“qS r   )Úname)r#   Úqr   r   r(   r)   C   s    z#Queues.__init__.<locals>.<dictcomp>)ÚsuperÚ__init__r   ÚaliasesÚdefault_exchangeÚdefault_routing_keyÚcreate_missingr   ÚautoexchangeÚmax_priorityr    r   r*   r
   ÚaddÚ
add_compat)	ÚselfÚqueuesr2   r4   r5   r6   r3   r-   r.   ©Ú	__class__r   r(   r0   8   s   
$€ÿzQueues.__init__c                    s,   z| j | W S  ty   tƒ  |¡ Y S w ©N)r1   ÚKeyErrorr/   Ú__getitem__©r9   r-   r;   r   r(   r?   H   s
   ÿzQueues.__getitem__c                    s<   | j r
|js
| j |_tƒ  ||¡ |jr|| j|j< d S d S r=   )r2   Úexchanger/   Ú__setitem__Úaliasr1   )r9   r-   Úqueuer;   r   r(   rB   N   s   ÿzQueues.__setitem__c                 C   s   | j r|  |  |¡¡S t|ƒ‚r=   )r4   r7   Únew_missingr>   r@   r   r   r(   Ú__missing__U   s   zQueues.__missing__c                 K   s&   t |tƒs| j|fi |¤ŽS |  |¡S )a¯  Add new queue.

        The first argument can either be a :class:`kombu.Queue` instance,
        or the name of a queue.  If the former the rest of the keyword
        arguments are ignored, and options are simply taken from the queue
        instance.

        Arguments:
            queue (kombu.Queue, str): Queue to add.
            exchange (kombu.Exchange, str):
                if queue is str, specifies exchange name.
            routing_key (str): if queue is str, specifies binding key.
            exchange_type (str): if queue is str, specifies type of exchange.
            **options (Any): Additional declaration options used when
                queue is a str.
        )r    r
   r8   Ú_add)r9   rD   Úkwargsr   r   r(   r7   Z   s   

z
Queues.addc                 K   s>   |  d| d¡¡ |d d u r||d< |  tj|fi |¤Ž¡S )NÚrouting_keyÚbinding_key)Ú
setdefaultÚgetrG   r
   Ú	from_dict)r9   r-   Úoptionsr   r   r(   r8   o   s   zQueues.add_compatc                 C   s`   |j d u s|j jdkr| j|_ |js| j|_| jd ur)|jd u r#i |_|  |j¡ || |j< |S )NÚ )rA   r-   r2   rI   r3   r6   Úqueue_argumentsÚ_set_max_priority)r9   rD   r   r   r(   rG   v   s   


zQueues._addc                 C   s*   d|vr| j d ur| d| j i¡S d S d S )Nzx-max-priority)r6   Úupdate)r9   Úargsr   r   r(   rQ   ‚   s   ÿzQueues._set_max_priorityr   c                 C   s\   | j }|sdS dd„ t| ¡ ƒD ƒ}|rtd |¡|ƒS |d d td |dd… ¡|ƒ S )z/Format routing table into string for log dumps.rO   c                 S   s   g | ]\}}t  ¡  |¡‘qS r   )ÚQUEUE_FORMATÚstripÚformat)r#   Ú_r.   r   r   r(   Ú
<listcomp>‹   s    ÿz!Queues.format.<locals>.<listcomp>Ú
r   r   N)Úconsume_fromÚsortedr*   Ú
textindentÚjoin)r9   r   Úindent_firstÚactiveÚinfor   r   r(   rV   †   s   
ÿ$zQueues.formatc                 K   s,   | j |fi |¤Ž}| jdur|| j|j< |S )z±Add new task queue that'll be consumed from.

        The queue will be active even when a subset has been selected
        using the :option:`celery worker -Q` option.
        N)r7   Ú_consume_fromr-   )r9   rD   rH   r.   r   r   r(   Ú
select_add‘   s   
zQueues.select_addc                    s$   |r‡ fdd„t |ƒD ƒˆ _dS dS )z¤Select a subset of currently defined queues to consume from.

        Arguments:
            include (Sequence[str], str): Names of queues to consume from.
        c                    ó   i | ]}|ˆ | “qS r   r   )r#   r-   ©r9   r   r(   r)   £   s    
ÿz!Queues.select.<locals>.<dictcomp>N)r   ra   )r9   Úincluder   rd   r(   Úselectœ   s
   
ÿÿzQueues.selectc                    sN   ˆ r#t ˆ ƒ‰ | jdu r|  ‡ fdd„| D ƒ¡S ˆ D ]}| j |d¡ qdS dS )z´Deselect queues so that they won't be consumed from.

        Arguments:
            exclude (Sequence[str], str): Names of queues to avoid
                consuming from.
        Nc                 3   s   � | ]	}|ˆ vr|V  qd S r=   r   )r#   r$   ©Úexcluder   r(   Ú	<genexpr>²   s   € z"Queues.deselect.<locals>.<genexpr>)r   ra   rf   Úpop)r9   rh   rD   r   rg   r(   Údeselect§   s   
ùzQueues.deselectc                 C   s   t ||  |¡|ƒS r=   )r
   r5   r@   r   r   r(   rE   ·   s   zQueues.new_missingc                 C   s   | j d ur| j S | S r=   )ra   rd   r   r   r(   rZ   º   s   
zQueues.consume_from)NNTNNN)r   T)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ra   r0   r?   rB   rF   r7   r8   rG   rQ   rV   rb   rf   rk   rE   ÚpropertyrZ   Ú__classcell__r   r   r;   r(   r   )   s*    þ
r   c                   @   sN  e Zd ZdZeZeZeZeZeZ	dZ
dZdZdZdZdd„ Zedd„ ƒZedd	„ ƒZ		d0d
d„Zd1dd„Zdd„ Zd1dd„Z									d2dd„Z							d3dd„Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZejd d„ ƒZed!d"„ ƒZed#d$„ ƒZejd%d$„ ƒZed&d'„ ƒZ e Z!ed(d)„ ƒZ"ed*d+„ ƒZ#ed,d-„ ƒZ$d.d/„ Z%dS )4r   zApp AMQP API: app.amqp.Ni   c                 C   s*   || _ | j| jdœ| _| j j | j¡ d S )N)r   é   )ÚappÚ
as_task_v1Ú
as_task_v2Útask_protocolsÚ_confÚbind_toÚ_handle_conf_update)r9   rs   r   r   r(   r0   á   s
   þzAMQP.__init__c                 C   ó   | j | jjj S r=   )rv   rs   ÚconfÚtask_protocolrd   r   r   r(   Úcreate_task_messageé   ó   zAMQP.create_task_messagec                 C   ó   |   ¡ S r=   )Ú_create_task_senderrd   r   r   r(   Úsend_task_messageí   ó   zAMQP.send_task_messagec                 C   sp   | j j}|j}|d u r|j}|d u r|j}|s$|jr$t|j| j|d�f}|d u r+| jn|}|  	|| j||||¡S )N)rA   rI   )
rs   r{   Útask_default_routing_keyÚtask_create_missing_queuesÚtask_queue_max_priorityÚtask_default_queuer
   r2   r5   Ú
queues_cls)r9   r:   r4   r5   r6   r{   r3   r   r   r(   r   ñ   s$   
þÿþzAMQP.Queuesc                 C   s&   t j| j|p| j| j d|¡| jd�S )zReturn the current task router.r„   )rs   )Ú_routesÚRouterr   r:   rs   Úeither)r9   r:   r4   r   r   r(   r‰     s   ÿþzAMQP.Routerc                 C   s   t  | jjj¡| _d S r=   )rˆ   Úpreparers   r{   Útask_routesÚ_rtablerd   r   r   r(   Úflush_routes  s   zAMQP.flush_routesc                 K   s:   |d u r	| j jj}| j|f||pt| jj ¡ ƒdœ|¤ŽS )N)Úacceptr:   )rs   r{   Úaccept_contentr   Úlistr:   rZ   Úvalues)r9   Úchannelr:   r�   Úkwr   r   r(   ÚTaskConsumer  s   
ÿþýzAMQP.TaskConsumerr   Fc           !         sú  |pd}|pi }t |ttfƒstdƒ‚t |tƒstdƒ‚|r<|  |d¡ |p*| j ¡ }|p0| jj}t	|t
|d� |d�}t |	tjƒr`|  |	d¡ |pN| j ¡ }|pT| jj}t	|t
|	d� |d�}	t |tƒsk|oj| ¡ }t |	tƒsv|	ou|	 ¡ }	|d u r€t|| jƒ}|d u rŠt|| jƒ}|sŽ|}‡ fdd	„|p–g D ƒ}i d
d“d|“d|“d|“d|“d|	“d|“d|“d|
“d||g“d|“d|“d|“d|“d|pËtƒ “d|“d|“||dœ¥} t| ||pÞddœ||||||dœf|rù|||||||
||	dœ	d �S d d �S )!Nr   ú!task args must be a list or tupleú(task keyword arguments must be a mappingÚ	countdown©Úseconds)ÚtzÚexpiresc                    rc   r   r   )r#   Úheader©rN   r   r(   r)   D  s    z#AMQP.as_task_v2.<locals>.<dictcomp>ÚlangÚpyÚtaskÚidÚshadowÚetaÚgroupÚgroup_indexÚretriesÚ	timelimitÚroot_idÚ	parent_idÚargsreprÚ
kwargsreprÚoriginÚignore_resultÚreplaced_task_nesting)Ústamped_headersÚstampsrO   ©Úcorrelation_idÚreply_to)Ú	callbacksÚerrbacksÚchainÚchord)	Úuuidr©   rª   r-   rS   rH   r§   r¤   rœ   r   )r    r‘   ÚtupleÚ	TypeErrorr   Ú_verify_secondsrs   ÚnowÚtimezoner   r   ÚnumbersÚRealÚstrÚ	isoformatr   Úargsrepr_maxsizeÚkwargsrepr_maxsizer   r   )!r9   Útask_idr-   rS   rH   r˜   r¤   Úgroup_idr¦   rœ   r§   r¸   rµ   r¶   r´   Ú
time_limitÚsoft_time_limitÚcreate_sent_eventr©   rª   r£   r·   r½   r¾   r­   r®   r«   r¬   r°   r¯   rN   r±   r   r   rž   r(   ru     sÀ   

ÿÿ

ÿþýüûúùø	÷
öõôóò
ñðïíþüÿö÷òèzAMQP.as_task_v2c                 K   s  |pd}|pi }| j }t|ttfƒstdƒ‚t|tƒstdƒ‚|r5|  |d¡ |p-| j ¡ }|t	|d� }t|	t
jƒrO|  |	d¡ |pG| j ¡ }|t	|	d� }	|oT| ¡ }|	oZ|	 ¡ }	ti ||paddœ|||||||
||	|||||f||d	œ|rˆ||t|ƒt|ƒ|
||	d
œd�S d d�S )Nr   r–   r—   r˜   r™   rœ   rO   r²   )r¡   r¢   rS   rH   r¥   r¦   r§   r¤   rœ   Úutcrµ   r¶   r¨   Útasksetr¸   )r¹   r-   rS   rH   r§   r¤   rœ   r   )rÊ   r    r‘   rº   r»   r   r¼   rs   r½   r   r¿   rÀ   rÂ   r   r   )r9   rÅ   r-   rS   rH   r˜   r¤   rÆ   r¦   rœ   r§   r¸   rµ   r¶   r´   rÇ   rÈ   rÉ   r©   rª   r£   r½   r¾   Úcompat_kwargsrÊ   r   r   r(   rt   v  sf   
þñøùéázAMQP.as_task_v1c                 C   s   |t k rt|› d|›�ƒ‚|S )Nz is out of range: )ÚINT_MINÚ
ValueError)r9   ÚsÚwhatr   r   r(   r¼   ²  s   zAMQP._verify_secondsc                    sÀ   | j jj‰| j jj‰| j jj‰| j‰| j‰tjj	‰tjj
‰tjj	‰tjj
‰ tjj	‰tjj
‰| j‰| j‰| j jj‰	| j jj‰
| j jj‰	 	 	 	 	 	 d‡ ‡‡‡‡‡‡‡‡‡	‡
‡‡‡‡‡fdd„	}|S )Nc                    sl  |d u rˆn|}|\}}}}|r|  |¡ |r|  |¡ |}|d u r(|d u r(ˆ}|d ur<t|tƒr9|ˆ| }}n|j}|
d u rTz|jj}
W n	 tyO   Y nw |
pSˆ}
|d u rjz|jj}W n tyi   d}Y nw |rn|sx|dkrxd|}}n|d u r‰|jjp�ˆ}|pˆ|jpˆˆ	}|d u r—|r—t|t	ƒs—|g}|d u r�ˆn|}|r©t
ˆfi |¤Žnˆ}ˆr¹ˆ||||||||d� | j|f|||	pÂˆ
|pÅˆ|||
||dœ	|¤Ž}ˆ rÛˆ|||||d� ˆ�rt|tƒrùˆ||d ||d |d |d	 |d
 d� nˆ||d ||d |d |d	 |d d� |�r4|�pˆ}|}t|tƒ�r!|j}|  |||dœ¡ |jd|| ||d� |S )NÚdirectrO   )Úsenderr   rA   rI   Údeclarer   r   Úretry_policy)	rA   rI   Ú
serializerÚcompressionÚretryrÔ   Údelivery_moderÓ   r   )rÒ   r   r   rA   rI   r¢   r   r   r¤   r¥   )rÒ   rÅ   r¡   rS   rH   r¤   rË   rS   rH   rË   )rD   rA   rI   z	task-sent)r×   rÔ   )rR   r    rÁ   r-   rA   rØ   ÚAttributeErrorÚtyperI   r   ÚdictÚpublishrº   r   )Úproducerr-   ÚmessagerA   rI   rD   Úevent_dispatcherr×   rÔ   rÕ   rØ   rÖ   rÓ   r   Úexchange_typerH   Úheaders2r   r   r   ÚqnameÚ_rpÚretÚevdÚexname©Úafter_receiversÚbefore_receiversÚdefault_compressorÚdefault_delivery_modeÚdefault_evdr2   Údefault_policyÚdefault_queueÚdefault_retryÚdefault_rkeyÚdefault_serializerr:   Úsend_after_publishÚsend_before_publishÚsend_task_sentÚsent_receiversr   r(   r�   Ì  s®   


ÿÿÿüÿø	÷ÿ

ý
ý
ýÿz3AMQP._create_task_sender.<locals>.send_task_message)NNNNNNNNNNNN)rs   r{   Útask_publish_retryÚtask_publish_retry_policyÚtask_default_delivery_moderî   r:   r   Úbefore_task_publishÚsendÚ	receiversÚafter_task_publishÚ	task_sentÚ_event_dispatcherr2   rƒ   Útask_serializerÚtask_compression)r9   r�   r   rç   r(   r€   ·  s0   





,úbzAMQP._create_task_senderc                 C   rz   r=   )r:   rs   r{   r†   rd   r   r   r(   rî   0  r~   zAMQP.default_queuec                 C   s   |   | jjj¡S )u"   Queue nameâ‡’ declaration mapping.)r   rs   r{   Útask_queuesrd   r   r   r(   r:   4  s   zAMQP.queuesc                 C   s
   |   |¡S r=   )r   )r9   r:   r   r   r(   r:   9  ó   
c                 C   s   | j d u r	|  ¡  | j S r=   )r�   rŽ   rd   r   r   r(   r   =  s   
zAMQP.routesc                 C   r   r=   )r‰   rd   r   r   r(   ÚrouterC  r‚   zAMQP.routerc                 C   s   |S r=   r   )r9   Úvaluer   r   r(   r  G  s   c                 C   s0   | j d u rtj| j ¡  | _ | jjj| j _| j S r=   )Ú_producer_poolr   Ú	producersrs   Úconnection_for_writeÚpoolÚlimitrd   r   r   r(   Úproducer_poolK  s   
ÿzAMQP.producer_poolc                 C   s   t | jjj| jjjƒS r=   )r   rs   r{   Útask_default_exchangeÚtask_default_exchange_typerd   r   r   r(   r2   T  s   
ÿzAMQP.default_exchangec                 C   s
   | j jjS r=   )rs   r{   Ú
enable_utcrd   r   r   r(   rÊ   Y  r  zAMQP.utcc                 C   s   | j jjdd�S )NF)Úenabled)rs   ÚeventsÚ
Dispatcherrd   r   r   r(   rþ   ]  s   zAMQP._event_dispatcherc                 O   s&   d|v sd|v r|   ¡  |  ¡ | _d S )NrŒ   )rŽ   r‰   r  )r9   rS   rH   r   r   r(   ry   c  s   
zAMQP._handle_conf_update)NNN)NN)NNNNNNNr   NNNNNNFNNNNNNNFNNNr   )NNNNNNNr   NNNNNNFNNNNN)&rl   rm   rn   ro   r   r   r	   ÚBrokerConnectionr   r‡   r�   r  r5   rÃ   rÄ   r0   r   r}   r�   r‰   rŽ   r•   ru   rt   r¼   r€   rî   r:   Úsetterrp   r   r  r
  Úpublisher_poolr2   rÊ   rþ   ry   r   r   r   r(   r   Á   s‚    


ÿ

	
ø^
ú<y









r   )r   )-ro   r¿   Úcollectionsr   Úcollections.abcr   Údatetimer   Úweakrefr   Úkombur   r   r   r	   r
   r   Úkombu.commonr   Úkombu.utils.functionalr   Úkombu.utils.objectsr   Úceleryr   Úcelery.utils.nodenamesr   Úcelery.utils.safereprr   Úcelery.utils.textr   r\   Úcelery.utils.timer   rO   r   rˆ   Ú__all__rÍ   rT   r   r,   rÛ   r   r   r   r   r   r(   Ú<module>   s4     ÿ
 