o
    wvXj‘‚  ã                   @   sø  d Z ddlmZmZ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 dd
lmZmZ ddlmZmZ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) ddl*m+Z+ ej,d dkr‹dndZ-dZ.dZ/dZ0dZ1dZ2ee3ƒZ4eddƒZ5eddƒZ6G d d!„ d!e7ƒZ8G d"d#„ d#e9ƒZ:G d$d%„ d%e;ƒZ<G d&d'„ d'e7ƒZ=G d(d)„ d)e7ƒZ>G d*d+„ d+e)j?ƒZ?G d,d-„ d-e7ƒZ@G d.d/„ d/e@e)jAƒZBG d0d1„ d1e)jCƒZCG d2d3„ d3e)jDƒZDdS )4zPVirtual transport implementation.

Emulates the AMQ API for non-AMQ transports.
é    )Úabsolute_importÚprint_functionÚunicode_literalsN)Úarray)ÚOrderedDictÚdefaultdictÚ
namedtuple)Úcount)ÚFinalize)Úsleep)Úqueue_declare_ok_t)ÚResourceErrorÚChannelError)ÚEmptyÚitemsÚ	monotonic)Ú
get_logger)Ústr_to_bytesÚbytes_to_str)Úemergency_dump_state)Ú	FairCycle©Úuuid)Úbaseé   )ÚSTANDARD_EXCHANGE_TYPESé   ÚHó   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,   úY/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/virtual/base.pyÚencodeC   s   zBase64.encodec                 C   s   t  t|ƒ¡S r&   )r'   Ú	b64decoder   r)   r,   r,   r-   ÚdecodeF   ó   zBase64.decodeN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r.   r0   r,   r,   r,   r-   r%   @   s    r%   c                   @   ó   e Zd ZdZdS )ÚNotEquivalentErrorzAEntity declaration is not equivalent to the previous declaration.N©r2   r3   r4   r5   r,   r,   r,   r-   r7   J   ó    r7   c                   @   r6   )ÚUndeliverableWarningz.The message could not be delivered to a queue.Nr8   r,   r,   r,   r-   r:   N   r9   r:   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*   r<   r,   r,   r-   Ú__init__n   s   zBrokerState.__init__c                 C   s"   | j  ¡  | j ¡  | j ¡  d S r&   )r<   Úclearr=   r?   ©r*   r,   r,   r-   rA   s   s   

zBrokerState.clearc                 C   s   |||f| j v S r&   )r=   )r*   r    r!   r"   r,   r,   r-   Úhas_bindingx   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"   rF   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ÚbindingrB   r,   r-   Ú
<listcomp>�   s    z5BrokerState.queue_bindings_delete.<locals>.<listcomp>)r?   rK   rH   )r*   r    r=   r,   rB   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=   )rL   rF   rB   r,   r-   Ú	<genexpr>’   s
   € ÿ
ÿz-BrokerState.queue_bindings.<locals>.<genexpr>)r?   ©r*   r    r,   rB   r-   Úqueue_bindings‘   s   
þzBrokerState.queue_bindingsr&   )r2   r3   r4   r5   r<   r=   r?   r@   rA   rC   rG   rJ   rO   rR   r,   r,   r,   r-   r;   R   s    

	r;   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>   Ú_dirtyrE   Ú
_quick_ackÚ__setitem__Ú_quick_appendr
   Úrestore_unacked_onceÚ_on_collect)r*   rU   rV   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.
        )rV   ÚlenrW   rY   ©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)rV   Úmaxr_   rW   rY   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)rY   Ú_flushr\   )r*   ÚmessageÚdelivery_tagr,   r,   r-   Úappend×   s   z
QoS.appendc                 C   s
   | j | S r&   )rW   ©r*   rg   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)rY   rW   rK   rH   )r*   ÚdirtyÚ	deliveredÚ	dirty_tagr,   r,   r-   re   à   s   ÿûz
QoS._flushc                 C   ó   |   |¡ dS )z8Acknowledge message and remove from transactional state.N)rZ   ri   r,   r,   r-   Úackë   s   zQoS.ackFc                 C   s$   |r| j  | j| ¡ |  |¡ dS )z4Remove from transactional state and requeue message.N)rU   Ú_restore_at_beginningrW   rZ   ©r*   rg   Ú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)	re   rW   rU   Ú_restoreÚpopitemrH   ÚBaseExceptionrh   rA   )r*   rm   ÚerrorsÚrestoreÚpop_messageÚ_rf   Ú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.
        NrX   )Úfile)ÚstderrT)r^   Úcancelre   Úsysr   rW   Úrestore_at_shutdownrU   Ú
do_restoreÚgetattrÚprintÚRESTORING_FMTÚformatr_   r}   ÚlistÚzipÚRESTORE_PANIC_FMTr   rX   )r*   r   ÚstateÚ
unrestoredrx   Úmessagesr,   r,   r-   r]   
  s4   
ÿÿ
õ
úzQoS.restore_unacked_oncec                 O   ó   dS )zõRestore any pending unackwnowledged 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_visible)  s   	zQoS.restore_visible)r   ©Fr&   )r2   r3   r4   r5   rV   rW   rY   r‚   r@   rb   rd   rh   rj   re   rp   rt   r}   r]   r‘   r,   r,   r,   r-   rS   ˜   s"    

	

rS   c                       s*   e Zd ZdZd‡ fdd„	Zdd„ Z‡  ZS )ÚMessagezMessage object.Nc                    sx   || _ |d }| d¡}|r| || d¡¡}tt| ƒjd|||d | d¡| d¡| d¡|| d¡d	d
œ	|¤Ž d S )NÚ
propertiesÚbodyÚbody_encodingrg   úcontent-typeúcontent-encodingÚheadersÚdelivery_infozutf-8)	r•   rU   rg   Úcontent_typeÚcontent_encodingr™   r”   rš   Ú
postencoder,   )Ú_rawrj   Údecode_bodyÚsuperr“   r@   )r*   ÚpayloadrU   r�   r”   r•   ©Ú	__class__r,   r-   r@   8  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”   rU   Úencode_bodyr•   rj   Údictr™   rK   r›   rœ   )r*   Úpropsr•   r{   r™   r,   r,   r-   ÚserializableJ  s   
ÿ
ûzMessage.serializabler&   )r2   r3   r4   r5   r@   r¨   Ú__classcell__r,   r,   r¢   r-   r“   5  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-   Ú_gete  ó   zAbstractChannel._getc                 C   r«   )zPut `message` onto `queue`.z$Virtual channels must implement _putr¬   )r*   r    rf   r,   r,   r-   Ú_puti  r°   zAbstractChannel._putc                 C   r«   )z!Remove all messages from `queue`.z&Virtual channels must implement _purger¬   rQ   r,   r,   r-   Ú_purgem  r°   zAbstractChannel._purgec                 C   rŽ   )z<Return the number of messages in `queue` as an :class:`int`.r   r,   rQ   r,   r,   r-   Ú_sizeq  s   zAbstractChannel._sizec                 O   ro   )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-   Ú_deleteu  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~  ó   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‡  r¸   zAbstractChannel._has_queuec                 C   s
   |  |¡S )z-Poll a list of queues for available messages.)rj   )r*   ÚcycleÚcallbackr®   r,   r,   r-   Ú_poll�  ó   
zAbstractChannel._pollc                 C   s   |   |¡}|||ƒ d S r&   )r¯   )r*   r    r»   rf   r,   r,   r-   Ú_get_and_deliver”  s   
z AbstractChannel._get_and_deliverr&   )r2   r3   r4   r5   r¯   r±   r²   r³   rµ   r·   r¹   r¼   r¾   r,   r,   r,   r-   rª   Z  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			d`dd„Zdadd„Zdbdd„Zdadd„Zdd„ Z		dcdd„Z		dcdd„Z		dddd„Z		dddd„Zd d!„ Zd"d#„ Z d$d%„ Z!d&d'„ Z"d(d)„ Z#d*d+„ Z$d,d-„ Z%ded.d/„Z&ded0d1„Z'ded2d3„Z(ded4d5„Z)		dfd6d7„Z*d8d9„ Z+d:d;„ Z,dgd<d=„Z-dhd>d?„Z.d@dA„ Z/dBdC„ Z0didDdE„Z1dFdG„ Z2		djdHdI„Z3dkdJdK„Z4dLdM„ Z5dhdNdO„Z6dhdPdQ„Z7dRdS„ Z8dTdU„ Z9dVdW„ Z:e;dXdY„ ƒZ<e;dZd[„ ƒZ=e;d\d]„ ƒZ>ded^d_„Z?dS )lÚ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„tˆ j	ƒD ƒˆ _	z	ˆ j j
 ¡ ˆ _W n tyB   td tˆ j jƒˆ j j¡dƒ‚w ˆ j jj}ˆ jD ]}z
tˆ ||| ƒ W qK ty`   Y qKw d S )NFc                    s   i | ]	\}}||ˆ ƒ“qS r,   r,   )rL   ÚtypÚclsrB   r,   r-   Ú
<dictcomp>Ñ  s    ÿz$Channel.__init__.<locals>.<dictcomp>z1No free channel ids, current={0}, channel_max={1})é   é
   )Ú
connectionr>   Ú
_consumersÚ_cycleÚ_tag_to_queueÚ_active_queuesÚ_qosÚclosedr   Úexchange_typesÚ_avail_channel_idsrK   Ú
channel_idÚ
IndexErrorr   r‡   r_   ÚchannelsÚchannel_maxÚclientÚtransport_optionsÚfrom_transport_optionsÚsetattrrH   )r*   rÇ   r�   ÚtoptsÚopt_namer,   rB   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 {0!r} in vhost {1!r}ú/©é2   rÆ   úChannel.exchange_declareÚ404N)ÚtypeÚdurableÚauto_deleter$   Útable)r‹   r<   r   r‡   rÇ   rÔ   Úvirtual_hostÚtypeofÚ
equivalentr7   ÚNOT_EQUIVALENT_FMTrH   )	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‹   r<   rK   )r*   r!   rì   rè   Úrkeyr{   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 {0!r} in vhost {1!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‹   rR   rå   Úprepare_bindrµ   rO   )	r*   r    rì   rí   r�   r!   r"   r�   Úmetar,   r,   r-   rï     s   
ÿzChannel.queue_deletec                 C   s   |   |¡ d S r&   )rï   rQ   r,   r,   r-   Úafter_reply_message_received!  r1   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‹   rC   rG   r<   rD   rå   rô   rh   Úsupports_fanoutÚ_queue_bind)r*   r    r!   r"   r$   r�   rã   rõ   r,   r,   r-   Ú
queue_bind,  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,   )rL   rõ   ©Úbinding_metar,   r-   rN   J  s    z(Channel.queue_unbind.<locals>.<listcomp>)r‹   rJ   rî   rH   rå   rô   )r*   r    r!   r"   r$   r�   rã   r,   r   r-   Úqueue_unbind=  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î   )rL   r!   rð   Úpatternr    rB   r,   r-   rP   M  s   € þþz(Channel.list_bindings.<locals>.<genexpr>©r‹   r<   rB   r,   rB   r-   Úlist_bindingsL  s   
ÿzChannel.list_bindingsc                 K   ó
   |   |¡S )z%Remove all ready messages from queue.r´   r¶   r,   r,   r-   Úqueue_purgeQ  r½   zChannel.queue_purgec                 C   s   t ƒ S r&   r   rB   r,   r,   r-   Ú_next_delivery_tagU  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*   rf   r!   r"   r�   r,   r,   r-   Úbasic_publishX  s   
ÿÿzChannel.basic_publishc                 C   sJ   |   |d | j¡\|d< }|d }|j||  ¡ d� |d j||d� d S )Nr•   r”   )r–   rg   rš   ©r!   r"   )r¥   r–   Úupdater  )r*   rf   r!   r"   r–   r§   r,   r,   r-   r	  b  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©rU   )r“   Úqosrh   rg   )Úraw_messagerf   ©r»   Úno_ackr*   r,   r-   Ú	_callbacku  s   z(Channel.basic_consume.<locals>._callbackN)rÊ   rË   rh   rÇ   Ú
_callbacksrÈ   rE   Ú_reset_cycle)r*   r    r  r»   Úconsumer_tagr�   r  r,   r  r-   Úbasic_consumep  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È   rI   r  rÊ   rK   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  rh   rg   r   )r*   r    r  r�   rf   r,   r,   r-   Ú	basic_getŒ  s   ÿzChannel.basic_getc                 C   s   | j  |¡ dS )zAcknowledge message.N)r  rp   )r*   rg   Ú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  r}   r­   )r*   rs   r,   r,   r-   Úbasic_recoverš  s   
zChannel.basic_recoverc                 C   s   | j j||d� dS )zReject message.©rs   N)r  rt   rr   r,   r,   r-   Úbasic_reject   s   zChannel.basic_rejectc                 C   s   || j _dS )zmChange QoS settings for this channel.

        Note:
            Only `prefetch_count` is supported.
        N)r  rV   )r*   Úprefetch_sizerV   Úapply_globalr,   r,   r-   Ú	basic_qos¤  s   zChannel.basic_qosc                 C   s   t | jjƒS r&   )rˆ   r‹   r<   rB   r,   r,   r-   Úget_exchanges­  ó   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‹   r<   rH   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:
            str: queue name -- must return the string `default`
                if no queues matched.
        Nr  )rÀ   rå   Úlookuprî   rH   ÚwarningsÚwarnr:   Ú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*   rf   rš   r    r,   r,   r-   ru   Ø  s   ÿþzChannel._restorec                 C   r  r&   )ru   )r*   rf   r,   r,   r-   rq   á  rk   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  rb   Ú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¡   rU   )Ú
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”   )rD   Údefault_priority)r*   r•   r5  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Ì   r]   rÉ   ÚcloserÇ   Úclose_channelrÎ   )r*   Úconsumerr,   r,   r-   r:    s   




zChannel.closec                 C   s"   |r| j  |¡ |¡|fS ||fS r&   )Úcodecsrj   r.   ©r*   r•   Úencodingr,   r,   r-   r¥     s   zChannel.encode_bodyc                 C   s   |r| j  |¡ |¡S |S r&   )r=  rj   r0   r>  r,   r,   r-   rŸ     s   zChannel.decode_bodyc                 C   s   t | j| jtƒ| _d S r&   )r   r¾   rË   r   rÉ   rB   r,   r,   r-   r  $  s   

ÿzChannel._reset_cyclec                 C   s   | S r&   r,   rB   r,   r,   r-   Ú	__enter__(  s   zChannel.__enter__c                 G   s   |   ¡  d S r&   )r:  )r*   Úexc_infor,   r,   r-   Ú__exit__+  r%  zChannel.__exit__c                 C   s   | j jS )z/Broker state containing exchanges and bindings.)rÇ   r‹   rB   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Ì   rS   rB   r,   r,   r-   r  3  s   
zChannel.qosc                 C   s   | j d u r	|  ¡  | j S r&   )rÉ   r  rB   r,   r,   r-   rº   :  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”   r5  )	rc   ÚminÚintÚmax_priorityÚmin_priorityÚ	TypeErrorr  rH   r6  )r*   rf   Úreverser5  r,   r,   r-   Ú_get_message_priority@  s   ÿý
ÿzChannel._get_message_priority)NrÚ   FFNFF)FF)NF)r÷   r÷   FN)Nr÷   Nr’   )r   r   F)rÚ   r&   )NN)NNNNN)T)@r2   r3   r4   r5   r“   rS   rƒ   r¦   r   rÎ   rý   r%   r=  r–   r	   Ú_delivery_tagsrÀ   rÖ   r6  rF  rE  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,  ru   rq   r2  r4  r7  r9  r:  r¥   rŸ   r  r@  rB  Úpropertyr‹   r  rº   rI  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 t| ƒ |¡ |j ¡ | _d S r&   )r    rL  r@   rÔ   rU   )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,   )rL   ÚqÚeÚrr,   r,   r-   rN   \  s    ÿz+Management.get_bindings.<locals>.<listcomp>)rU   r  rB   r,   r,   r-   Úget_bindings[  s   ÿzManagement.get_bindingsc                 C   s   | j  ¡  d S r&   )rU   r:  rB   r,   r,   r-   r:  _  r1   zManagement.close)r2   r3   r4   r5   r@   rQ  r:  r©   r,   r,   r¢   r-   rL  T  s
    rL  c                   @   s¶   e Zd ZdZe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 ) Ú	TransportznVirtual transport.

    Arguments:
        client (kombu.Connection): The client this is a transport for.
    Ng      ð?iÿÿ  FrÚ   Útopic)ÚasynchronousÚexchange_typeÚ
heartbeatsc                 K   s`   || _ g | _g | _i | _|  | j| jt¡| _|j 	d¡}|d ur#|| _
ttt| jddƒƒ| _d S )NÚpolling_intervalr   éÿÿÿÿ)rÔ   rÒ   Ú_avail_channelsr  ÚCycleÚ_drain_channelr   rº   rÕ   rj   rW  r   ÚARRAY_TYPE_HÚrangerÓ   rÏ   )r*   rÔ   r�   rW  r,   r,   r-   r@   Š  s   
ÿzTransport.__init__c                 C   s:   z| j  ¡ W S  ty   |  |¡}| j |¡ | Y S w r&   )rY  rK   rÑ   r¿   rÒ   rh   )r*   rÇ   rU   r,   r,   r-   Úcreate_channel—  s   
ýzTransport.create_channelc                 C   sT   z%| j  |j¡ z| j |¡ W n	 ty   Y nw W d |_d S W d |_d S d |_w r&   )rÏ   rh   rÐ   rÒ   rI   r  rÇ   )r*   rU   r,   r,   r-   r;  Ÿ  s   ÿÿ
þzTransport.close_channelc                 C   s   | j  |  | ¡¡ | S r&   )rY  rh   r^  rB   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º   r:  rY  rÒ   rK   ÚLookupError)r*   rÇ   Ú	chan_listrU   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º   rj   rW  r0  r   Úsocketr®   r   )r*   rÇ   r®   Ú
time_startrj   rW  r,   r,   r-   r2  »  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: {0})rH   r‡   r  ÚloggerÚwarningÚW_NO_CONSUMERSÚ_reject_inbound_message)r*   rf   r    r»   r,   r,   r-   r0  Ì  s   ÿÿþzTransport._deliverc                 C   sH   | j D ]}|r!|j||d�}|j ||j¡ |j|jdd�  d S qd S )Nr  Tr  )rÒ   r“   r  rh   rg   r   )r*   r  rU   rf   r,   r,   r-   rh  Ù  s   
üÿz!Transport._reject_inbound_messagec                 C   s0   |r|| j vrtd ||¡ƒ‚| j | |ƒ d S )Nz.Message for queue {0!r} without consumers: {1})r  rH   r‡   )r*   rU   rf   r    r,   r,   r-   Úon_message_readyá  s   ÿÿzTransport.on_message_readyc                 C   s   |j ||d�S )N)r»   r®   )r2  )r*   rU   r»   r®   r,   r,   r-   r[  è  r1   zTransport._drain_channelc                 C   s   | j ddœS )NÚ	localhost)ÚportÚhostname)Údefault_portrB   r,   r,   r-   Údefault_connection_paramsë  s   z#Transport.default_connection_paramsr&   )!r2   r3   r4   r5   r¿   r   rZ  rL  r;   r‹   rº   rm  rÒ   r  rW  rÓ   r   rR  Ú
implementsÚextendÚ	frozensetr@   r^  r;  r_  rb  r2  r0  rh  ri  r[  rK  rn  r,   r,   r,   r-   rR  c  s:    
ý


rR  )Er5   Ú
__future__r   r   r   r'   rc  r�   r(  r   Úcollectionsr   r   r   Ú	itertoolsr	   Úmultiprocessing.utilr
   Útimer   Úamqp.protocolr   Úkombu.exceptionsr   r   Ú
kombu.fiver   r   r   Ú	kombu.logr   Úkombu.utils.encodingr   r   Úkombu.utils.divr   Úkombu.utils.schedulingr   Úkombu.utils.uuidr   Úkombu.transportr   r!   r   Úversion_infor\  r*  rç   rg  r†   rŠ   r2   re  r   r#   Úobjectr%   Ú	Exceptionr7   ÚUserWarningr:   r;   rS   r“   rª   Ú
StdChannelr¿   rL  rR  r,   r,   r,   r-   Ú<module>   sX    


F %?   >