o
    —¨Êhö…  ã                   @  sî  d Z ddlmZ ddlZddlZddlZddl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mZ dd
lmZ ddl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$ ddl%m&Z& ddl'm(Z( ddl)m*Z* erŠddl+m,Z, dZ-dZ.dZ/dZ0dZ1dZ2ee3ƒZ4eddƒZ5eddƒZ6G d d!„ d!ƒZ7G d"d#„ d#e8ƒZ9G d$d%„ d%e:ƒZ;G d&d'„ d'ƒZ<G d(d)„ d)ƒZ=G d*d+„ d+ej>ƒZ>G d,d-„ d-ƒZ?G d.d/„ d/e?ej@ƒZAG d0d1„ d1ejBƒZBG d2d3„ d3ejCƒZCdS )4zPVirtual transport implementation.

Emulates the AMQ API for non-AMQ transports.
é    )ÚannotationsN)Úarray)ÚOrderedDictÚdefaultdictÚ
namedtuple)Úcount)ÚFinalize)ÚEmpty)Ú	monotonicÚsleep)ÚTYPE_CHECKING)Úqueue_declare_ok_t)ÚChannelErrorÚResourceError)Ú
get_logger)Úbase)Úemergency_dump_state)Úbytes_to_strÚstr_to_bytes)Ú	FairCycle©Úuuidé   )ÚSTANDARD_EXCHANGE_TYPES)ÚTracebackTypeÚHzlMessage could not be delivered: No queues bound to exchange {exchange!r} using binding key {routing_key!r}.
zkCannot redeclare exchange {0!r} in vhost {1!r} with different type, durable, autodelete or arguments value.z;Requeuing undeliverable message for queue %r: No consumers.z)Restoring {0!r} unacknowledged message(s)z#UNABLE TO RESTORE {0} MESSAGES: {1}Úbinding_key_t)ÚqueueÚexchangeÚrouting_keyÚqueue_binding_t)r   r   Ú	argumentsc                   @  s    e Zd ZdZdd„ Zdd„ ZdS )ÚBase64zBase64 codec.c                 C  s   t t t|ƒ¡ƒS ©N)r   Úbase64Ú	b64encoder   ©ÚselfÚs© r)   úN/var/www/html/env/lib/python3.10/site-packages/kombu/transport/virtual/base.pyÚencodeF   s   zBase64.encodec                 C  s   t  t|ƒ¡S r#   )r$   Ú	b64decoder   r&   r)   r)   r*   ÚdecodeI   ó   zBase64.decodeN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r+   r-   r)   r)   r)   r*   r"   C   s    r"   c                   @  ó   e Zd ZdZdS )ÚNotEquivalentErrorzAEntity declaration is not equivalent to the previous declaration.N©r/   r0   r1   r2   r)   r)   r)   r*   r4   M   ó    r4   c                   @  r3   )ÚUndeliverableWarningz.The message could not be delivered to a queue.Nr5   r)   r)   r)   r*   r7   Q   r6   r7   c                   @  sV   e Zd Zd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S )ÚBrokerStatez2Broker state holds exchanges, queues and bindings.Nc                 C  s&   |d u ri n|| _ i | _ttƒ| _d S r#   )Ú	exchangesÚbindingsr   ÚsetÚqueue_index)r'   r9   r)   r)   r*   Ú__init__r   s   zBrokerState.__init__c                 C  s"   | j  ¡  | j ¡  | j ¡  d S r#   )r9   Úclearr:   r<   ©r'   r)   r)   r*   r>   w   s   

zBrokerState.clearc                 C  s   |||f| j v S r#   )r:   )r'   r   r   r   r)   r)   r*   Úhas_binding|   s   zBrokerState.has_bindingc                 C  s.   t |||ƒ}| j ||¡ | j|  |¡ d S r#   )r   r:   Ú
setdefaultr<   Úadd)r'   r   r   r   r!   Úkeyr)   r)   r*   Úbinding_declare   s   zBrokerState.binding_declarec                 C  sB   t |||ƒ}z| j|= W n
 ty   Y d S w | j|  |¡ d S r#   )r   r:   ÚKeyErrorr<   Úremove)r'   r   r   r   rC   r)   r)   r*   Úbinding_delete„   s   ÿzBrokerState.binding_deletec                   s<   zˆ j  |¡}W n
 ty   Y d S w ‡ fdd„|D ƒ d S )Nc                   s   g | ]	}ˆ j  |d ¡‘qS r#   )r:   Úpop)Ú.0Úbindingr?   r)   r*   Ú
<listcomp>“   s    z5BrokerState.queue_bindings_delete.<locals>.<listcomp>)r<   rH   rE   )r'   r   r:   r)   r?   r*   Úqueue_bindings_delete�   s   ÿz!BrokerState.queue_bindings_deletec                   s   ‡ fdd„ˆ j | D ƒS )Nc                 3  s&   � | ]}t |j|jˆ j| ƒV  qd S r#   )r    r   r   r:   )rI   rC   r?   r)   r*   Ú	<genexpr>–   s
   € ÿ
ÿz-BrokerState.queue_bindings.<locals>.<genexpr>)r<   ©r'   r   r)   r?   r*   Úqueue_bindings•   s   
þzBrokerState.queue_bindingsr#   )r/   r0   r1   r2   r9   r:   r<   r=   r>   r@   rD   rG   rL   rO   r)   r)   r)   r*   r8   U   s    

	r8   c                   @  s~   e Zd ZdZd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d„Zdd„ Zddd„Zdd„ ZdS )ÚQoSzêQuality of Service guarantees.

    Only supports `prefetch_count` at this point.

    Arguments:
    ---------
        channel (ChannelT): Connection channel.
        prefetch_count (int): Initial prefetch count (defaults to 0).
    r   NTc                 C  sR   || _ |pd| _tƒ | _d| j_tƒ | _| jj| _| jj	| _
t| | jdd�| _d S )Nr   Fr   )Úexitpriority)ÚchannelÚprefetch_countr   Ú
_deliveredÚrestoredr;   Ú_dirtyrB   Ú
_quick_ackÚ__setitem__Ú_quick_appendr   Úrestore_unacked_onceÚ_on_collect)r'   rR   rS   r)   r)   r*   r=   ·   s   


ÿzQoS.__init__c                 C  s$   | j }| pt| jƒt| jƒ |k S )z�Return true if the channel can be consumed from.

        Used to ensure the client adhers to currently active
        prefetch limits.
        )rS   ÚlenrT   rV   ©r'   Úpcountr)   r)   r*   Úcan_consumeÆ   s   zQoS.can_consumec                 C  s,   | j }|rt|t| jƒt| jƒ  dƒS dS )a’  Return the maximum number of messages allowed to be returned.

        Returns an estimated number of messages that a consumer may be allowed
        to consume at once from the broker.  This is used for services where
        bulk 'get message' calls are preferred to many individual 'get message'
        calls - like SQS.

        Returns
        -------
            int: greater than zero.
        r   N)rS   Úmaxr\   rT   rV   r]   r)   r)   r*   Úcan_consume_max_estimateÏ   s   ÿzQoS.can_consume_max_estimatec                 C  s   | j r|  ¡  |  ||¡ dS )z&Append message to transactional state.N)rV   Ú_flushrY   )r'   ÚmessageÚdelivery_tagr)   r)   r*   Úappendß   s   z
QoS.appendc                 C  s
   | j | S r#   )rT   ©r'   rd   r)   r)   r*   Úgetå   ó   
zQoS.getc                 C  s>   | j }| j}	 z| ¡ }W n
 ty   Y dS w | |d¡ q)z'Flush dirty (acked/rejected) tags from.r   N)rV   rT   rH   rE   )r'   ÚdirtyÚ	deliveredÚ	dirty_tagr)   r)   r*   rb   è   s   ÿûz
QoS._flushc                 C  ó   |   |¡ dS )z8Acknowledge message and remove from transactional state.N)rW   rf   r)   r)   r*   Úackó   s   zQoS.ackFc                 C  s$   |r| j  | j| ¡ |  |¡ dS )z4Remove from transactional state and requeue message.N)rR   Ú_restore_at_beginningrT   rW   ©r'   rd   Úrequeuer)   r)   r*   Úreject÷   s   z
QoS.rejectc              
   C  s–   |   ¡  | j}g }| jj}|j}|rEz|ƒ \}}W n	 ty"   Y n#w z||ƒ W n tyB } z| ||f¡ W Y d}~nd}~ww |s| ¡  |S )z$Restore all unacknowledged messages.N)	rb   rT   rR   Ú_restoreÚpopitemrE   ÚBaseExceptionre   r>   )r'   rj   ÚerrorsÚrestoreÚpop_messageÚ_rc   Úexcr)   r)   r*   Úrestore_unackedý   s(   ÿ€ÿø
zQoS.restore_unackedc                 C  sÞ   | j  ¡  |  ¡  |du rtjn|}| j}| jr| jjsdS t	|ddƒr*|r(J ‚dS z@|r_t
t t| jƒ¡|d� |  ¡ }|rett|Ž ƒ\}}t
t t|ƒ|¡|d� t||d� W d|_dS W d|_dS W d|_dS d|_w )zÅRestore all unacknowledged messages at shutdown/gc collect.

        Note:
        ----
            Can only be called once for each instance, subsequent
            calls will be ignored.
        NrU   )Úfile)ÚstderrT)r[   Úcancelrb   Úsysr|   rT   Úrestore_at_shutdownrR   Ú
do_restoreÚgetattrÚprintÚRESTORING_FMTÚformatr\   rz   ÚlistÚzipÚRESTORE_PANIC_FMTr   rU   )r'   r|   ÚstateÚ
unrestoredru   Úmessagesr)   r)   r*   rZ     s4   
ÿÿ
õ
úzQoS.restore_unacked_oncec                 O  ó   dS )a  Restore any pending unacknowledged messages.

        To be filled in for visibility_timeout style implementations.

        Note:
        ----
            This is implementation optional, and currently only
            used by the Redis transport.
        Nr)   )r'   ÚargsÚkwargsr)   r)   r*   Úrestore_visible2  ó    zQoS.restore_visible)r   ©Fr#   )r/   r0   r1   r2   rS   rT   rV   r   r=   r_   ra   re   rg   rb   rm   rq   rz   rZ   rŽ   r)   r)   r)   r*   rP   œ   s"    
	

 rP   c                      s*   e Zd ZdZd‡ fdd„	Zdd„ Z‡  ZS )ÚMessagezMessage object.Nc                   st   || _ |d }| d¡}|r| || d¡¡}tƒ jd|||d | d¡| d¡| d¡|| d¡d	d
œ	|¤Ž d S )NÚ
propertiesÚbodyÚbody_encodingrd   úcontent-typeúcontent-encodingÚheadersÚdelivery_infoúutf-8)	r“   rR   rd   Úcontent_typeÚcontent_encodingr—   r’   r˜   Ú
postencoder)   )Ú_rawrg   Údecode_bodyÚsuperr=   )r'   ÚpayloadrR   r�   r’   r“   ©Ú	__class__r)   r*   r=   A  s$   
÷

özMessage.__init__c                 C  sJ   | j }| j | j| d¡¡\}}t| jƒ}| dd ¡ ||| j| j	|dœS )Nr”   Úcompression)r“   r’   r•   r–   r—   )
r’   rR   Úencode_bodyr“   rg   Údictr—   rH   rš   r›   )r'   Úpropsr“   rx   r—   r)   r)   r*   ÚserializableS  s   
ÿ
ûzMessage.serializabler#   )r/   r0   r1   r2   r=   r§   Ú__classcell__r)   r)   r¡   r*   r‘   >  s    r‘   c                   @  s\   e 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d„Zdd„ ZdS )ÚAbstractChannelzþAbstract channel interface.

    This is an abstract class defining the channel methods
    you'd usually want to implement in a virtual channel.

    Note:
    ----
        Do not subclass directly, but rather inherit
        from :class:`Channel`.
    Nc                 C  ó   t dƒ‚)zGet next message from `queue`.z$Virtual channels must implement _get©ÚNotImplementedError)r'   r   Útimeoutr)   r)   r*   Ú_geto  ó   zAbstractChannel._getc                 C  rª   )zPut `message` onto `queue`.z$Virtual channels must implement _putr«   )r'   r   rc   r)   r)   r*   Ú_puts  r¯   zAbstractChannel._putc                 C  rª   )z!Remove all messages from `queue`.z&Virtual channels must implement _purger«   rN   r)   r)   r*   Ú_purgew  r¯   zAbstractChannel._purgec                 C  r‹   )z<Return the number of messages in `queue` as an :class:`int`.r   r)   rN   r)   r)   r*   Ú_size{  s   zAbstractChannel._sizec                 O  rl   )z�Delete `queue`.

        Note:
        ----
            This just purges the queue, if you need to do more you can
            override this method.
        N©r±   )r'   r   rŒ   r�   r)   r)   r*   Ú_delete  s   zAbstractChannel._deletec                 K  r‹   )z´Create new queue.

        Note:
        ----
            Your transport can override this method if it needs
            to do something whenever a new queue is declared.
        Nr)   ©r'   r   r�   r)   r)   r*   Ú
_new_queue‰  r�   zAbstractChannel._new_queuec                 K  r‹   )z²Verify that queue exists.

        Returns
        -------
            bool: Should return :const:`True` if the queue exists
                or :const:`False` otherwise.
        Tr)   rµ   r)   r)   r*   Ú
_has_queue’  s   zAbstractChannel._has_queuec                 C  s
   |  |¡S )z-Poll a list of queues for available messages.)rg   )r'   ÚcycleÚcallbackr­   r)   r)   r*   Ú_pollœ  ó   
zAbstractChannel._pollc                 C  s   |   |¡}|||ƒ d S r#   )r®   )r'   r   r¹   rc   r)   r)   r*   Ú_get_and_deliver   s   
z AbstractChannel._get_and_deliverr#   )r/   r0   r1   r2   r®   r°   r±   r²   r´   r¶   r·   rº   r¼   r)   r)   r)   r*   r©   c  s    

	

r©   c                   @  s   e Zd ZdZeZeZdZeeƒZ	dZ
deƒ iZdZedƒZdZdZdZdZd	Zd
d„ Z			djdd„Zdkdd„Zdldd„Zdkdd„Zdd„ Z		dmdd„Z		dmdd„Z		dndd„Z		dndd„Zd d!„ Zd"d#„ Z d$d%„ Z!d&d'„ Z"d(d)„ Z#d*d+„ Z$d,d-„ Z%dod.d/„Z&dod0d1„Z'dod2d3„Z(dod4d5„Z)		dpd6d7„Z*d8d9„ Z+d:d;„ Z,dqd<d=„Z-drd>d?„Z.d@dA„ Z/dBdC„ Z0dsdDdE„Z1dFdG„ Z2		dtdHdI„Z3dudJdK„Z4dLdM„ Z5drdNdO„Z6drdPdQ„Z7dRdS„ Z8dTdU„ Z9dvd^d_„Z:e;d`da„ ƒZ<e;dbdc„ ƒZ=e;ddde„ ƒZ>dodfdg„Z?dhdi„ Z@dS )wÚChannelz‘Virtual channel.

    Arguments:
    ---------
        connection (ConnectionT): The transport instance this
            channel is part of.
    TFr$   r   N)r”   Údeadletter_queuer   é	   c              	     s�   |ˆ _ tƒ ˆ _d ˆ _i ˆ _g ˆ _d ˆ _dˆ _‡ fdd„ˆ j 	¡ D ƒˆ _ˆ  
¡ ˆ _ˆ j jj}ˆ jD ]}z
tˆ ||| ƒ W q0 tyE   Y q0w d S )NFc                   s   i | ]	\}}||ˆ ƒ“qS r)   r)   )rI   ÚtypÚclsr?   r)   r*   Ú
<dictcomp>Þ  s    ÿz$Channel.__init__.<locals>.<dictcomp>)Ú
connectionr;   Ú
_consumersÚ_cycleÚ_tag_to_queueÚ_active_queuesÚ_qosÚclosedÚexchange_typesÚitemsÚ_get_free_channel_idÚ
channel_idÚclientÚtransport_optionsÚfrom_transport_optionsÚsetattrrE   )r'   rÃ   r�   ÚtoptsÚopt_namer)   r?   r*   r=   Ô  s&   
ÿ


ÿýzChannel.__init__Údirectc           	   	   C  sÀ   |pd}|p	d| }|r$|| j jvr"td || jjjpd¡dddƒ‚dS z#| j j| }|  |¡ ||||||¡sEt	t
 || jjjpBd¡ƒ‚W dS  ty_   ||||pTi g d	œ| j j|< Y dS w )
zDeclare exchange.rÔ   zamq.%sz*NOT_FOUND - no exchange {!r} in vhost {!r}ú/©é2   é
   úChannel.exchange_declareÚ404N)ÚtypeÚdurableÚauto_deleter!   Útable)rˆ   r9   r   r„   rÃ   rÎ   Úvirtual_hostÚtypeofÚ
equivalentr4   ÚNOT_EQUIVALENT_FMTrE   )	r'   r   rÛ   rÜ   rÝ   r!   ÚnowaitÚpassiveÚprevr)   r)   r*   Úexchange_declareë  s:   ÿýþÿýûÿrÙ   c                 C  s:   |   |¡D ]\}}}| j|ddd� q| jj |d¡ dS )z'Delete `exchange` and all its bindings.T)Ú	if_unusedÚif_emptyN)Ú	get_tableÚqueue_deleterˆ   r9   rH   )r'   r   rç   rã   Úrkeyrx   r   r)   r)   r*   Úexchange_delete	  s   zChannel.exchange_deletec                 K  sh   |pdt ƒ  }|r"| j|fi |¤Žs"td || jjjpd¡dddƒ‚| j|fi |¤Ž t||  	|¡dƒS )zDeclare queue.z
amq.gen-%sz'NOT_FOUND - no queue {!r} in vhost {!r}rÕ   rÖ   úChannel.queue_declarerÚ   r   )
r   r·   r   r„   rÃ   rÎ   rß   r¶   r   r²   )r'   r   rä   r�   r)   r)   r*   Úqueue_declare  s   ÿýrí   c           	      K  sj   |r	|   |¡r	dS | j |¡D ]\}}}|  |¡ ||||¡}| j||g|¢R i |¤Ž q| j |¡ dS )zDelete queue.N)r²   rˆ   rO   rà   Úprepare_bindr´   rL   )	r'   r   rç   rè   r�   r   r   rŒ   Úmetar)   r)   r*   rê     s   
ÿzChannel.queue_deletec                 C  s   |   |¡ d S r#   )rê   rN   r)   r)   r*   Úafter_reply_message_received'  r.   z$Channel.after_reply_message_receivedÚ c                 C  rª   )Nz(transport does not support exchange_bindr«   ©r'   ÚdestinationÚsourcer   rã   r!   r)   r)   r*   Úexchange_bind*  r¯   zChannel.exchange_bindc                 C  rª   )Nz*transport does not support exchange_unbindr«   ró   r)   r)   r*   Úexchange_unbind.  r¯   zChannel.exchange_unbindc                 K  s‚   |pd}| j  |||¡rdS | j  ||||¡ | j j|  dg ¡}|  |¡ ||||¡}| |¡ | jr?| j	|g|¢R Ž  dS dS )z.Bind `queue` to `exchange` with `routing key`.z
amq.directNrÞ   )
rˆ   r@   rD   r9   rA   rà   rï   re   Úsupports_fanoutÚ_queue_bind)r'   r   r   r   r!   r�   rÞ   rð   r)   r)   r*   Ú
queue_bind2  s   
ÿ
ÿzChannel.queue_bindc                   sh   | j  |||¡ z|  |¡}W n
 ty   Y d S w |  |¡ ||||¡‰ ‡ fdd„|D ƒ|d d …< d S )Nc                   s   g | ]}|ˆ kr|‘qS r)   r)   )rI   rð   ©Úbinding_metar)   r*   rK   P  s    z(Channel.queue_unbind.<locals>.<listcomp>)rˆ   rG   ré   rE   rà   rï   )r'   r   r   r   r!   r�   rÞ   r)   rû   r*   Úqueue_unbindC  s   ÿ
ÿzChannel.queue_unbindc                   s   ‡ fdd„ˆ j jD ƒS )Nc                 3  s0   � | ]}ˆ   |¡D ]\}}}|||fV  q	qd S r#   )ré   )rI   r   rë   Úpatternr   r?   r)   r*   rM   S  s   € þþz(Channel.list_bindings.<locals>.<genexpr>©rˆ   r9   r?   r)   r?   r*   Úlist_bindingsR  s   
ÿzChannel.list_bindingsc                 K  ó
   |   |¡S )z%Remove all ready messages from queue.r³   rµ   r)   r)   r*   Úqueue_purgeW  r»   zChannel.queue_purgec                 C  s   t ƒ S r#   r   r?   r)   r)   r*   Ú_next_delivery_tag[  s   zChannel._next_delivery_tagc                 K  sB   |   |||¡ |r|  |¡j|||fi |¤ŽS | j||fi |¤ŽS )zPublish message.)Ú_inplace_augment_messagerà   Údeliverr°   )r'   rc   r   r   r�   r)   r)   r*   Úbasic_publish^  s   
ÿÿzChannel.basic_publishc                 C  sJ   |   |d | j¡\|d< }|d }|j||  ¡ d� |d j||d� d S )Nr“   r’   )r”   rd   r˜   ©r   r   )r¤   r”   Úupdater  )r'   rc   r   r   r”   r¦   r)   r)   r*   r  h  s   
ÿþ
þz Channel._inplace_augment_messagec                   sJ   |ˆj |< ˆj |¡ ‡ ‡‡fdd„}|ˆjj|< ˆj |¡ ˆ ¡  dS )zConsume from `queue`.c                   s*   ˆj | ˆd�}ˆsˆj ||j¡ ˆ |ƒS )N©rR   )r‘   Úqosre   rd   )Úraw_messagerc   ©r¹   Úno_ackr'   r)   r*   Ú	_callback{  s   z(Channel.basic_consume.<locals>._callbackN)rÆ   rÇ   re   rÃ   Ú
_callbacksrÄ   rB   Ú_reset_cycle)r'   r   r  r¹   Úconsumer_tagr�   r  r)   r  r*   Úbasic_consumev  s   
zChannel.basic_consumec                 C  sh   || j v r2| j  |¡ |  ¡  | j |d¡}z| j |¡ W n	 ty'   Y nw | jj |d¡ dS dS )z Cancel consumer by consumer tag.N)	rÄ   rF   r  rÆ   rH   rÇ   Ú
ValueErrorrÃ   r  )r'   r  r   r)   r)   r*   Úbasic_cancel†  s   
ÿøzChannel.basic_cancelc                 K  sD   z| j |  |¡| d�}|s| j ||j¡ |W S  ty!   Y dS w )z+Get message by direct access (synchronous).r	  N)r‘   r®   r
  re   rd   r	   )r'   r   r  r�   rc   r)   r)   r*   Ú	basic_get’  s   ÿzChannel.basic_getc                 C  s   | j  |¡ dS )zAcknowledge message.N)r
  rm   )r'   rd   Úmultipler)   r)   r*   Ú	basic_ackœ  ó   zChannel.basic_ackc                 C  s   |r| j  ¡ S tdƒ‚)zRecover unacked messages.z'Does not support recover(requeue=False))r
  rz   r¬   )r'   rp   r)   r)   r*   Úbasic_recover   s   
zChannel.basic_recoverc                 C  s   | j j||d� dS )zReject message.©rp   N)r
  rq   ro   r)   r)   r*   Úbasic_reject¦  s   zChannel.basic_rejectc                 C  s   || j _dS )zzChange QoS settings for this channel.

        Note:
        ----
            Only `prefetch_count` is supported.
        N)r
  rS   )r'   Úprefetch_sizerS   Úapply_globalr)   r)   r*   Ú	basic_qosª  s   zChannel.basic_qosc                 C  s   t | jjƒS r#   )r…   rˆ   r9   r?   r)   r)   r*   Úget_exchanges´  s   zChannel.get_exchangesc                 C  s   | j j| d S )z%Get table of bindings for `exchange`.rÞ   rÿ   )r'   r   r)   r)   r*   ré   ·  r  zChannel.get_tablec                 C  s6   z
| j j| d }W n ty   |}Y nw | j| S )z.Get the exchange type instance for `exchange`.rÛ   )rˆ   r9   rE   rÊ   )r'   r   ÚdefaultrÛ   r)   r)   r*   rà   »  s   ÿ
zChannel.typeofc                 C  sŒ   |du r| j }|s|p|gS z|  |¡ |  |¡|||¡}W n ty)   g }Y nw |sD|durDt ttj	||d�ƒ¡ |  
|¡ |g}|S )záFind all queues matching `routing_key` for the given `exchange`.

        Returns
        -------
            list[str]: queue names -- must return `[default]`
                if default is set and no queues matched.
        Nr  )r¾   rà   Úlookupré   rE   ÚwarningsÚwarnr7   ÚUNDELIVERABLE_FMTr„   r¶   )r'   r   r   r   ÚRr)   r)   r*   Ú_lookupÃ  s&   

þÿ

ÿ
zChannel._lookupc                 C  s@   |j }| ¡ }d|d< |  |d |d ¡D ]}|  ||¡ qdS )z.Redeliver message to its original destination.TÚredeliveredr   r   N)r˜   r§   r&  r°   )r'   rc   r˜   r   r)   r)   r*   rr   à  s   þýzChannel._restorec                 C  r  r#   )rr   )r'   rc   r)   r)   r*   rn   ê  rh   zChannel._restore_at_beginningc                 C  sN   |p| j j}| jr$| j ¡ r$t| dƒr| j| j|d�S | j| j	||d�S t
ƒ ‚)NÚ	_get_many©r­   )rÃ   Ú_deliverrÄ   r
  r_   Úhasattrr(  rÇ   rº   r¸   r	   )r'   r­   r¹   r)   r)   r*   Údrain_eventsí  s   
zChannel.drain_eventsc                 C  s   t || jƒs| j|| d�S |S )z1Convert raw message to :class:`Message` instance.)r    rR   )Ú
isinstancer‘   )r'   r  r)   r)   r*   Úmessage_to_pythonõ  s   zChannel.message_to_pythonc                 C  s>   |pi }|  di ¡ |  d|p| j¡ ||||pi |pi dœS )zPrepare message data.r˜   Úpriority)r“   r–   r•   r—   r’   )rA   Údefault_priority)r'   r“   r/  rš   r›   r—   r’   r)   r)   r*   Úprepare_messageû  s   üzChannel.prepare_messagec                 C  rª   )z´Enable/disable message flow.

        Raises
        ------
            NotImplementedError: as flow
                is not implemented by the base virtual implementation.
        z%virtual channels do not support flow.r«   )r'   Úactiver)   r)   r*   Úflow  s   zChannel.flowc                 C  sp   | j s3d| _ t| jƒD ]}|  |¡ q| jr| j ¡  | jdur(| j ¡  d| _| jdur3| j 	| ¡ d| _
dS )zTClose channel.

        Cancel all consumers, and requeue unacked messages.
        TN)rÉ   r…   rÄ   r  rÈ   rZ   rÅ   ÚcloserÃ   Úclose_channelrÊ   )r'   Úconsumerr)   r)   r*   r4    s   




zChannel.closec                 C  s.   |r|  ¡ dkr| j |¡ |¡|fS ||fS ©Nr™   )ÚlowerÚcodecsrg   r+   ©r'   r“   Úencodingr)   r)   r*   r¤   $  s   zChannel.encode_bodyc                 C  s&   |r|  ¡ dkr| j |¡ |¡S |S r7  )r8  r9  rg   r-   r:  r)   r)   r*   rž   )  s   zChannel.decode_bodyc                 C  s   t | j| jtƒ| _d S r#   )r   r¼   rÇ   r	   rÅ   r?   r)   r)   r*   r  .  s   

ÿzChannel._reset_cyclec                 C  s   | S r#   r)   r?   r)   r)   r*   Ú	__enter__2  s   zChannel.__enter__Úexc_typeútype[BaseException] | NoneÚexc_valúBaseException | NoneÚexc_tbúTracebackType | NoneÚreturnÚNonec                 C  s   |   ¡  d S r#   )r4  )r'   r=  r?  rA  r)   r)   r*   Ú__exit__5  s   zChannel.__exit__c                 C  s   | j jS )z/Broker state containing exchanges and bindings.)rÃ   rˆ   r?   r)   r)   r*   rˆ   =  s   zChannel.statec                 C  s   | j du r|  | ¡| _ | j S )z&:class:`QoS` manager for this channel.N)rÈ   rP   r?   r)   r)   r*   r
  B  s   
zChannel.qosc                 C  s   | j d u r	|  ¡  | j S r#   )rÅ   r  r?   r)   r)   r*   r¸   I  s   
zChannel.cyclec              
   C  sV   zt tt|d d ƒ| jƒ| jƒ}W n tttfy!   | j}Y nw |r)| j| S |S )z©Get priority from message.

        The value is limited to within a boundary of 0 to 9.

        Note:
        ----
            Higher value has more priority.
        r’   r/  )	r`   ÚminÚintÚmax_priorityÚmin_priorityÚ	TypeErrorr  rE   r0  )r'   rc   Úreverser/  r)   r)   r*   Ú_get_message_priorityO  s   	ÿý
ÿzChannel._get_message_priorityc                 C  s`   t | jjƒ}td| jjd ƒD ]}||vr | jj |¡ |  S qtd t| jj	ƒ| jj¡dƒ‚)Nr   z/No free channel ids, current={}, channel_max={})é   rØ   )
r;   rÃ   Ú_used_channel_idsÚrangeÚchannel_maxre   r   r„   r\   Úchannels)r'   Úused_channel_idsrÍ   r)   r)   r*   rÌ   c  s   þ
þýzChannel._get_free_channel_id)NrÔ   FFNFF)FF)NF)rò   rò   FN)Nrò   Nr�   )r   r   F)rÔ   r#   )NN)NNNNN)T)r=  r>  r?  r@  rA  rB  rC  rD  )Ar/   r0   r1   r2   r‘   rP   r€   r¥   r   rÊ   rø   r"   r9  r”   r   Ú_delivery_tagsr¾   rÐ   r0  rI  rH  r=   ræ   rì   rî   rê   rñ   rö   r÷   rú   rý   r   r  r  r  r  r  r  r  r  r  r  r  r  ré   rà   r&  rr   rn   r,  r.  r1  r3  r4  r¤   rž   r  r<  rE  Úpropertyrˆ   r
  r¸   rL  rÌ   r)   r)   r)   r*   r½   ¥  s˜    	

þ



ÿ
ÿ
ÿ
ÿ






ÿ





ÿ








r½   c                      s0   e Zd ZdZ‡ fdd„Zdd„ Zdd„ Z‡  ZS )Ú
Managementz'Base class for the AMQP management API.c                   s   t ƒ  |¡ |j ¡ | _d S r#   )rŸ   r=   rÎ   rR   )r'   Ú	transportr¡   r)   r*   r=   w  s   zManagement.__init__c                 C  s   dd„ | j  ¡ D ƒS )Nc                 S  s   g | ]\}}}|||d œ‘qS ))rô   rõ   r   r)   )rI   ÚqÚeÚrr)   r)   r*   rK   |  s    ÿz+Management.get_bindings.<locals>.<listcomp>)rR   r   r?   r)   r)   r*   Úget_bindings{  s   ÿzManagement.get_bindingsc                 C  s   | j  ¡  d S r#   )rR   r4  r?   r)   r)   r*   r4    r.   zManagement.close)r/   r0   r1   r2   r=   rZ  r4  r¨   r)   r)   r¡   r*   rU  t  s
    rU  c                   @  s°   e Zd ZdZeZeZeZdZdZ	dZ
dZdZdZejjjdeddgƒ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d„Zedd„ ƒZdS ) Ú	Transportz|Virtual transport.

    Arguments:
    ---------
        client (kombu.Connection): The client this is a transport for.
    Ng      ð?iÿÿ  FrÔ   Útopic)ÚasynchronousÚexchange_typeÚ
heartbeatsc                 K  s\   || _ tƒ | _g | _g | _i | _|  | j| jt¡| _	|j
 d¡}|d ur'|| _ttƒ| _d S )NÚpolling_interval)rÎ   r8   rˆ   rQ  Ú_avail_channelsr  ÚCycleÚ_drain_channelr	   r¸   rÏ   rg   r`  r   ÚARRAY_TYPE_HrN  )r'   rÎ   r�   r`  r)   r)   r*   r=   ¨  s   zTransport.__init__c                 C  s:   z| j  ¡ W S  ty   |  |¡}| j |¡ | Y S w r#   )ra  rH   Ú
IndexErrorr½   rQ  re   )r'   rÃ   rR   r)   r)   r*   Úcreate_channelµ  s   
ýzTransport.create_channelc                 C  sl   z1z	| j  |j¡ W n	 ty   Y nw z| j |¡ W n	 ty%   Y nw W d |_d S W d |_d S d |_w r#   )rN  rF   rÍ   r  rQ  rÃ   )r'   rR   r)   r)   r*   r5  ½  s   þÿÿ
þzTransport.close_channelc                 C  s   | j  |  | ¡¡ | S r#   )ra  re   rf  r?   r)   r)   r*   Úestablish_connectionË  s   zTransport.establish_connectionc              	   C  sP   | j  ¡  | j| jfD ]}|r%z| ¡ }W n	 ty   Y nw | ¡  |sqd S r#   )r¸   r4  ra  rQ  rH   ÚLookupError)r'   rÃ   Ú	chan_listrR   r)   r)   r*   Úclose_connectionÒ  s   
ÿú€ÿzTransport.close_connectionc                 C  s‚   t ƒ }| jj}| j}|r|r||kr|}	 z
|| j|d� W d S  ty?   |d ur5t ƒ | |kr5t ¡ ‚|d ur=t|ƒ Y nw q)Nr   r)  )	r
   r¸   rg   r`  r*  r	   Úsocketr­   r   )r'   rÃ   r­   Ú
time_startrg   r`  r)   r)   r*   r,  Ý  s"   ú€üýzTransport.drain_eventsc                 C  sX   |s	t d |¡ƒ‚z| j| }W n t y%   t t|¡ |  |¡ Y d S w ||ƒ d S )Nz.Received message without destination queue: {})rE   r„   r  ÚloggerÚwarningÚW_NO_CONSUMERSÚ_reject_inbound_message)r'   rc   r   r¹   r)   r)   r*   r*  î  s   ÿÿþzTransport._deliverc                 C  sH   | j D ]}|r!|j||d�}|j ||j¡ |j|jdd�  d S qd S )Nr	  Tr  )rQ  r‘   r
  re   rd   r  )r'   r  rR   rc   r)   r)   r*   rp  û  s   
üÿz!Transport._reject_inbound_messagec                 C  s0   |r|| j vrtd ||¡ƒ‚| j | |ƒ d S )Nz,Message for queue {!r} without consumers: {})r  rE   r„   )r'   rR   rc   r   r)   r)   r*   Úon_message_ready  s   ÿÿzTransport.on_message_readyc                 C  s   |j ||d�S )N)r¹   r­   )r,  )r'   rR   r¹   r­   r)   r)   r*   rc  
  r.   zTransport._drain_channelc                 C  s   | j ddœS )NÚ	localhost)ÚportÚhostname)Údefault_portr?   r)   r)   r*   Údefault_connection_params  s   z#Transport.default_connection_paramsr#   )r/   r0   r1   r2   r½   r   rb  rU  r¸   ru  rQ  r  r`  rP  r   r[  Ú
implementsÚextendÚ	frozensetr=   rf  r5  rg  rj  r,  r*  rp  rq  rc  rT  rv  r)   r)   r)   r*   r[  ƒ  s8    
ý

r[  )Dr2   Ú
__future__r   r$   rk  r~   r"  r   Úcollectionsr   r   r   Ú	itertoolsr   Úmultiprocessing.utilr   r   r	   Útimer
   r   Útypingr   Úamqp.protocolr   Úkombu.exceptionsr   r   Ú	kombu.logr   Úkombu.transportr   Úkombu.utils.divr   Úkombu.utils.encodingr   r   Úkombu.utils.schedulingr   Úkombu.utils.uuidr   r   r   Útypesr   rd  r$  râ   ro  rƒ   r‡   r/   rm  r   r    r"   Ú	Exceptionr4   ÚUserWarningr7   r8   rP   r‘   r©   Ú
StdChannelr½   rU  r[  r)   r)   r)   r*   Ú<module>   s^    


G #%B   R