o
    wvXjÂ_  ã                   @   sp  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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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+ zddl,m-Z- W n e.y�   ddlm-Z- Y nw dZ/dZ0e o™edƒ Z1dZ2eddƒZ3d"dd„Z4G dd„ de5ƒZ6G d d!„ d!e7ƒZ8dS )#z/Sending/Receiving Messages (Kombu integration).é    )Úabsolute_importÚunicode_literalsN)Ú
namedtuple)Ú	timedelta)ÚWeakValueDictionary)Ú
ConnectionÚConsumerÚExchangeÚProducerÚQueueÚpools)Ú	Broadcast)Ú
maybe_list)Úcached_property)Úsignals)ÚPY3ÚitemsÚstring_t)Ú
try_import)Úanon_nodename)Úsaferepr)Úindent)Úmaybe_make_awareé   )Úroutes)ÚMapping)ÚAMQPÚQueuesÚtask_messagei   €Ú
simplejsonzS
.> {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„t | ƒD ƒS )Nc                    s*   i | ]\}}t |tƒr| ˆ ¡n||“qS © )Ú
isinstanceÚbytesÚdecode)Ú.0ÚkÚv©Úencodingr&   úL/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/app/amqp.pyÚ
<dictcomp>2   s    ÿzutf8dict.<locals>.<dictcomp>)r   )Údr.   r&   r-   r/   Úutf8dict1   s   
ÿr2   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„ 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d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`.
        ha_policy (Sequence, str): Default HA policy for queues with none set.
        max_priority (int): Default x-max-priority for queues with none set.
    NTc           
      C   s¢   t  | ¡ tƒ | _|| _|| _|| _|| _|d u rtn|| _	|| _
|d ur1t|tƒs1dd„ |D ƒ}t|p5i ƒD ]\}}	t|	tƒrE|  |	¡n| j|fi |	¤Ž q7d S )Nc                 S   s   i | ]}|j |“qS r&   )Úname)r*   Úqr&   r&   r/   r0   R   ó    z#Queues.__init__.<locals>.<dictcomp>)ÚdictÚ__init__r   ÚaliasesÚdefault_exchangeÚdefault_routing_keyÚcreate_missingÚ	ha_policyr	   ÚautoexchangeÚmax_priorityr'   r   r   r   ÚaddÚ
add_compat)
ÚselfÚqueuesr9   r;   r<   r=   r>   r:   r3   r4   r&   r&   r/   r7   F   s   
$€ÿzQueues.__init__c                 C   s,   z| j | W S  ty   t | |¡ Y S w ©N)r8   ÚKeyErrorr6   Ú__getitem__©rA   r3   r&   r&   r/   rE   V   s
   ÿzQueues.__getitem__c                 C   s<   | j r
|js
| j |_t | ||¡ |jr|| j|j< d S d S rC   )r9   Úexchanger6   Ú__setitem__Úaliasr8   )rA   r3   Úqueuer&   r&   r/   rH   \   s   ÿzQueues.__setitem__c                 C   s   | j r|  |  |¡¡S t|ƒ‚rC   )r;   r?   Únew_missingrD   rF   r&   r&   r/   Ú__missing__c   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   r@   Ú_add)rA   rJ   Úkwargsr&   r&   r/   r?   h   s   

z
Queues.addc                 K   s>   |  d| d¡¡ |d d u r||d< |  tj|fi |¤Ž¡S )NÚrouting_keyÚbinding_key)Ú
setdefaultÚgetrM   r   Ú	from_dict)rA   r3   Úoptionsr&   r&   r/   r@   }   s   zQueues.add_compatc                 C   s‚   |j d u s|j jdkr| j|_ |js| j|_| jr'|jd u r!i |_|  |j¡ | jd ur:|jd u r4i |_|  	|j¡ || |j< |S )NÚ )
rG   r3   r9   rO   r:   r<   Úqueue_argumentsÚ_set_ha_policyr>   Ú_set_max_priority)rA   rJ   r&   r&   r/   rM   „   s   



zQueues._addc                 C   s4   | j }t|ttfƒr| dt|ƒdœ¡S ||d< d S )NÚnodes)úha-modez	ha-paramsrZ   )r<   r'   ÚlistÚtupleÚupdate)rA   ÚargsÚpolicyr&   r&   r/   rW   ”   s   ÿzQueues._set_ha_policyc                 C   s*   d|vr| j d ur| d| j i¡S d S d S )Nzx-max-priority)r>   r]   )rA   r^   r&   r&   r/   rX   ›   s   ÿzQueues._set_max_priorityr   c                 C   s\   | j }|sdS dd„ tt|ƒƒD ƒ}|rtd |¡|ƒS |d d td |dd… ¡|ƒ S )z/Format routing table into string for log dumps.rU   c                 S   s   g | ]\}}t  ¡  |¡‘qS r&   )ÚQUEUE_FORMATÚstripÚformat)r*   Ú_r4   r&   r&   r/   Ú
<listcomp>¤   s    ÿz!Queues.format.<locals>.<listcomp>Ú
r   r   N)Úconsume_fromÚsortedr   Ú
textindentÚjoin)rA   r   Úindent_firstÚactiveÚinfor&   r&   r/   rb   Ÿ   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)r?   Ú_consume_fromr3   )rA   rJ   rN   r4   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                    s   i | ]}|ˆ | “qS r&   r&   )r*   r3   ©rA   r&   r/   r0   ¼   s    
ÿz!Queues.select.<locals>.<dictcomp>N)r   rm   )rA   Úincluder&   ro   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 rC   r&   )r*   r+   ©Úexcluder&   r/   Ú	<genexpr>Ë   s   € z"Queues.deselect.<locals>.<genexpr>)r   rm   rq   Úpop)rA   rs   rJ   r&   rr   r/   ÚdeselectÀ   s   
ùzQueues.deselectc                 C   s   t ||  |¡|ƒS rC   )r   r=   rF   r&   r&   r/   rK   Ð   s   zQueues.new_missingc                 C   s   | j d ur| j S | S rC   )rm   ro   r&   r&   r/   rf   Ó   s   
zQueues.consume_from)NNTNNNN)r   T)Ú__name__Ú
__module__Ú__qualname__Ú__doc__rm   r7   rE   rH   rL   r?   r@   rM   rW   rX   rb   rn   rq   rv   rK   Úpropertyrf   r&   r&   r&   r/   r   6   s,    
þ
r   c                   @   sL  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)rA   r}   r&   r&   r/   r7   ú   s
   þzAMQP.__init__c                 C   ó   | j | jjj S rC   )r€   r}   ÚconfÚtask_protocolro   r&   r&   r/   Úcreate_task_message  ó   zAMQP.create_task_messagec                 C   ó   |   ¡ S rC   )Ú_create_task_senderro   r&   r&   r/   Úsend_task_message  ó   zAMQP.send_task_messagec              	   C   s€   | j j}|j}|d u r|j}|d u r|j}|d u r|j}|s+|jr+t|j| j|d�f}|d u r2| j	n|}|  
|| j|||||¡S )N)rG   rO   )r}   r…   Útask_default_routing_keyÚtask_create_missing_queuesÚtask_queue_ha_policyÚtask_queue_max_priorityÚtask_default_queuer   r9   r=   Ú
queues_cls)rA   rB   r;   r<   r=   r>   r…   r:   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Ž   )r}   )Ú_routesÚRouterr   rB   r}   Úeither)rA   rB   r;   r&   r&   r/   r”   !  s   ÿþzAMQP.Routerc                 C   s   t  | jjj¡| _d S rC   )r“   Úpreparer}   r…   Útask_routesÚ_rtablero   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)ÚacceptrB   )r}   r…   Úaccept_contentr   r[   rB   rf   Úvalues)rA   ÚchannelrB   rš   Úkwr&   r&   r/   ÚTaskConsumer*  s   
ÿþýzAMQP.TaskConsumerr   Fc                 C   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ƒ}tr¤|r•dd	„ |D ƒ}|ržd
d	„ |D ƒ}|
r¤t|
ƒ}
|s¨|}td|||||||	||g|||||p¼tƒ dœ||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                 S   ó   g | ]}t |ƒ‘qS r&   ©r2   ©r*   Úcallbackr&   r&   r/   rd   \  r5   z#AMQP.as_task_v2.<locals>.<listcomp>c                 S   r§   r&   r¨   ©r*   Úerrbackr&   r&   r/   rd   ^  r5   Úpy)ÚlangÚtaskÚidÚshadowÚetar¦   ÚgroupÚretriesÚ	timelimitÚroot_idÚ	parent_idÚargsreprÚ
kwargsreprÚoriginrU   ©Úcorrelation_idÚreply_to)Ú	callbacksÚerrbacksÚchainÚchord)	Úuuidr¶   r·   r3   r^   rN   r´   r²   r¦   r    )r'   r[   r\   Ú	TypeErrorr   Ú_verify_secondsr}   ÚnowÚtimezoner   r   ÚnumbersÚRealr   Ú	isoformatr   Úargsrepr_maxsizeÚkwargsrepr_maxsizeÚJSON_NEEDS_UNICODE_KEYSr2   r   r   )rA   Útask_idr3   r^   rN   r¢   r²   Úgroup_idr¦   r´   rÁ   r¾   r¿   r½   Ú
time_limitÚsoft_time_limitÚcreate_sent_eventr¶   r·   r±   rÀ   rÅ   rÆ   rº   r¸   r¹   r&   r&   r/   r   3  sœ   
ÿÿ

òþüÿö÷ã'ÙzAMQP.as_task_v2c                 K   sJ  |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ru|rfdd„ |D ƒ}|rod	d„ |D ƒ}|
rut|
ƒ}
ti ||p{d
dœ||||||	|||||||f||
dœ|r¡||t|ƒt|ƒ|	||dœd�S d d�S )Nr&   r    r¡   r¢   r£   r¦   c                 S   r§   r&   r¨   r©   r&   r&   r/   rd   «  r5   z#AMQP.as_task_v1.<locals>.<listcomp>c                 S   r§   r&   r¨   r«   r&   r&   r/   rd   ­  r5   rU   r»   )r¯   r°   r^   rN   r³   r´   r²   r¦   Úutcr¾   r¿   rµ   ÚtasksetrÁ   )rÂ   r3   r^   rN   r´   r²   r¦   r    )rÒ   r'   r[   r\   rÃ   r   rÄ   r}   rÅ   r   rÇ   rÈ   rÉ   rÌ   r2   r   r   )rA   rÍ   r3   r^   rN   r¢   r²   rÎ   r¦   r´   rÁ   r¾   r¿   r½   rÏ   rÐ   rÑ   r¶   r·   r±   rÅ   rÆ   Úcompat_kwargsrÒ   r&   r&   r/   r~   �  sr   
þòøùêâzAMQP.as_task_v1c                 C   s   |t k rtd||f ƒ‚|S )Nz%s is out of range: %r)ÚINT_MINÚ
ValueError)rA   Ú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ÚdirectrU   )Úsenderr#   rG   rO   Údeclarer!   r"   Úretry_policy)	rG   rO   Ú
serializerÚcompressionÚretryrÜ   Údelivery_moderÛ   r!   )rÚ   r#   r!   rG   rO   r°   r   r   r²   r³   )rÚ   rÍ   r¯   r^   rN   r²   rÓ   r^   rN   rÓ   )rJ   rG   rO   z	task-sent)rß   rÜ   )r]   r'   r   r3   rG   rà   ÚAttributeErrorÚtyperO   r   r6   Úpublishr\   r	   )Úproducerr3   ÚmessagerG   rO   rJ   Úevent_dispatcherrß   rÜ   rÝ   rà   rÞ   rÛ   r!   Úexchange_typerN   Úheaders2r"   r#   r$   ÚqnameÚ_rpÚretÚevdÚexname©Úafter_receiversÚbefore_receiversÚdefault_compressorÚdefault_delivery_modeÚdefault_evdr9   Údefault_policyÚdefault_queueÚdefault_retryÚdefault_rkeyÚdefault_serializerrB   Úsend_after_publishÚsend_before_publishÚsend_task_sentÚsent_receiversr&   r/   r‹   ì  s®   


ÿÿÿüÿø	÷ÿ

ý
ý
ýÿz3AMQP._create_task_sender.<locals>.send_task_message)NNNNNNNNNNNN)r}   r…   Útask_publish_retryÚtask_publish_retry_policyÚtask_default_delivery_moderõ   rB   r   Úbefore_task_publishÚsendÚ	receiversÚafter_task_publishÚ	task_sentÚ_event_dispatcherr9   r�   Útask_serializerÚresult_compression)rA   r‹   r&   rî   r/   rŠ   ×  s0   





,úbzAMQP._create_task_senderc                 C   r„   rC   )rB   r}   r…   r‘   ro   r&   r&   r/   rõ   P  rˆ   zAMQP.default_queuec                 C   s   |   | jjj¡S )u"   Queue nameâ‡’ declaration mapping.)r   r}   r…   Útask_queuesro   r&   r&   r/   rB   T  s   zAMQP.queuesc                 C   s
   |   |¡S rC   )r   )rA   rB   r&   r&   r/   rB   Y  ó   
c                 C   s   | j d u r	|  ¡  | j S rC   )r˜   r™   ro   r&   r&   r/   r   ]  s   
zAMQP.routesc                 C   r‰   rC   )r”   ro   r&   r&   r/   Úrouterc  rŒ   zAMQP.routerc                 C   s   |S rC   r&   )rA   Úvaluer&   r&   r/   r
  g  s   c                 C   s0   | j d u rtj| j ¡  | _ | jjj| j _| j S rC   )Ú_producer_poolr   Ú	producersr}   Úconnection_for_writeÚpoolÚlimitro   r&   r&   r/   Úproducer_poolk  s   
ÿzAMQP.producer_poolc                 C   s   t | jjj| jjjƒS rC   )r	   r}   r…   Útask_default_exchangeÚtask_default_exchange_typero   r&   r&   r/   r9   t  s   
ÿzAMQP.default_exchangec                 C   s
   | j jjS rC   )r}   r…   Ú
enable_utcro   r&   r&   r/   rÒ   y  r	  zAMQP.utcc                 C   s   | j jjdd�S )NF)Úenabled)r}   ÚeventsÚ
Dispatcherro   r&   r&   r/   r  }  s   zAMQP._event_dispatcherc                 O   s&   d|v sd|v r|   ¡  |  ¡ | _d S )Nr—   )r™   r”   r
  )rA   r^   rN   r&   r&   r/   rƒ   ƒ  s   
zAMQP._handle_conf_update)NNNN)NN)NNNNNNr   NNNNNNFNNNNNNNNN)NNNNNNr   NNNNNNFNNNNN)&rw   rx   ry   rz   r   r   r
   ÚBrokerConnectionr   r’   r˜   r  r=   rÊ   rË   r7   r   r‡   r‹   r”   r™   rŸ   r   r~   rÄ   rŠ   rõ   rB   Úsetterr{   r   r
  r  Úpublisher_poolr9   rÒ   r  rƒ   r&   r&   r&   r/   r   Ú   s€    


ÿ

	
ù\
úCy









r   )r%   )9rz   Ú
__future__r   r   rÇ   Úcollectionsr   Údatetimer   Úweakrefr   Úkombur   r   r	   r
   r   r   Úkombu.commonr   Úkombu.utils.functionalr   Úkombu.utils.objectsr   Úceleryr   Úcelery.fiver   r   r   Úcelery.localr   Úcelery.utils.nodenamesr   Úcelery.utils.safereprr   Úcelery.utils.textr   rh   Úcelery.utils.timer   rU   r   r“   Úcollections.abcr   ÚImportErrorÚ__all__rÕ   rÌ   r`   r   r2   r6   r   Úobjectr   r&   r&   r&   r/   Ú<module>   sD    þÿ
 %