o
    wvXjÁ  ã                   @   s  d Z ddlmZmZ ddlZddl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 dZdefdefdefdefdefdœZdd„ Zdd„ Zdd„ ZG dd„ deƒZG dd„ deƒZG dd„ deƒZede g d¢ƒdd �Z!G d!d"„ d"eƒZ"dS )#zBase transport interface.é    )Úabsolute_importÚunicode_literalsN)ÚRecoverableConnectionError)ÚChannelErrorÚConnectionError)Úitems)ÚMessage)Ú
dictfilter)Úcached_property)Úmaybe_s_to_ms)r   Ú
StdChannelÚ
ManagementÚ	Transportz	x-expireszx-message-ttlzx-max-lengthzx-max-length-byteszx-max-priority)ÚexpiresÚmessage_ttlÚ
max_lengthÚmax_length_bytesÚmax_priorityc                 K   s2   t tdd„ t|ƒD ƒƒƒ}|rt| fi |¤ŽS | S )a  Convert queue arguments to RabbitMQ queue arguments.

    This is the implementation for Channel.prepare_queue_arguments
    for AMQP-based transports.  It's used by both the pyamqp and librabbitmq
    transports.

    Arguments:
        arguments (Mapping):
            User-supplied arguments (``Queue.queue_arguments``).

    Keyword Arguments:
        expires (float): Queue expiry time in seconds.
            This will be converted to ``x-expires`` in int milliseconds.
        message_ttl (float): Message TTL in seconds.
            This will be converted to ``x-message-ttl`` in int milliseconds.
        max_length (int): Max queue length (in number of messages).
            This will be converted to ``x-max-length`` int.
        max_length_bytes (int): Max queue size in bytes.
            This will be converted to ``x-max-length-bytes`` int.
        max_priority (int): Max priority steps for queue.
            This will be converted to ``x-max-priority`` int.

    Returns:
        Dict: RabbitMQ compatible queue arguments.
    c                 s   s   � | ]
\}}t ||ƒV  qd S ©N)Ú_to_rabbitmq_queue_argument)Ú.0ÚkeyÚvalue© r   úQ/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/base.pyÚ	<genexpr>8   s
   € ÿ
ÿz.to_rabbitmq_queue_arguments.<locals>.<genexpr>)r	   Údictr   )Ú	argumentsÚoptionsÚpreparedr   r   r   Úto_rabbitmq_queue_arguments   s   

þr    c                 C   s&   t |  \}}||d ur||ƒfS |fS r   )ÚRABBITMQ_QUEUE_ARGUMENTS)r   r   ÚoptÚtypr   r   r   r   ?   s   r   c                 C   s   t d | j|¡ƒS )Nz<Transport {0.__module__}.{0.__name__} does not implement {1})ÚNotImplementedErrorÚformatÚ	__class__)ÚobjÚmethodr   r   r   Ú
_LeftBlankE   s
   ÿÿr)   c                   @   sL   e Zd ZdZdZdd„ Zdd„ Zdd„ Zd	d
„ Zdd„ Z	dd„ Z
dd„ ZdS )r   zStandard channel base class.Nc                 O   ó"   ddl m} || g|¢R i |¤ŽS )Nr   )ÚConsumer)Úkombu.messagingr+   )ÚselfÚargsÚkwargsr+   r   r   r   r+   P   ó   zStdChannel.Consumerc                 O   r*   )Nr   )ÚProducer)r,   r1   )r-   r.   r/   r1   r   r   r   r1   T   r0   zStdChannel.Producerc                 C   ó
   t | dƒ‚©NÚget_bindings©r)   ©r-   r   r   r   r4   X   ó   
zStdChannel.get_bindingsc                 C   ó   dS )z·Callback called after RPC reply received.

        Notes:
           Reply queue semantics: can be used to delete the queue
           after transient reply message received.
        Nr   )r-   Úqueuer   r   r   Úafter_reply_message_received[   s   z'StdChannel.after_reply_message_receivedc                 K   s   |S r   r   )r-   r   r/   r   r   r   Úprepare_queue_argumentsd   ó   z"StdChannel.prepare_queue_argumentsc                 C   s   | S r   r   r6   r   r   r   Ú	__enter__g   r<   zStdChannel.__enter__c                 G   s   |   ¡  d S r   )Úclose)r-   Úexc_infor   r   r   Ú__exit__j   ó   zStdChannel.__exit__)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Úno_ack_consumersr+   r1   r4   r:   r;   r=   r@   r   r   r   r   r   K   s    	r   c                   @   s    e Zd ZdZdd„ Zdd„ ZdS )r   z!AMQP Management API (incomplete).c                 C   ó
   || _ d S r   )Ú	transport)r-   rH   r   r   r   Ú__init__q   r7   zManagement.__init__c                 C   r2   r3   r5   r6   r   r   r   r4   t   r7   zManagement.get_bindingsN)rB   rC   rD   rE   rI   r4   r   r   r   r   r   n   s    r   c                   @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	Ú
Implementsz/Helper class used to define transport features.c                 C   s"   z| | W S  t y   t|ƒ‚w r   )ÚKeyErrorÚAttributeError)r-   r   r   r   r   Ú__getattr__{   s
   
ÿzImplements.__getattr__c                 C   s   || |< d S r   r   )r-   r   r   r   r   r   Ú__setattr__�   rA   zImplements.__setattr__c                 K   s   | j | fi |¤ŽS r   )r&   )r-   r/   r   r   r   Úextend„   s   zImplements.extendN)rB   rC   rD   rE   rM   rN   rO   r   r   r   r   rJ   x   s
    rJ   F)ÚdirectÚtopicÚfanoutÚheaders)ÚasynchronousÚexchange_typeÚ
heartbeatsc                   @   s  e Zd ZdZeZdZdZdZefZ	e
fZdZdZdZe ¡ 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dd„ Zejej e!j"e!j#ffdd„Z$d d!„ Z%d"d#„ Z&e'd$d%„ ƒZ(d&d'„ Z)e*d(d)„ ƒZ+e'd*d+„ ƒZ,e'd,d-„ ƒZ-dS )/r   zBase class for transports.NFúN/Ac                 K   rG   r   )Úclient)r-   rX   r/   r   r   r   rI   °   r7   zTransport.__init__c                 C   r2   )NÚestablish_connectionr5   r6   r   r   r   rY   ³   r7   zTransport.establish_connectionc                 C   r2   )NÚclose_connectionr5   ©r-   Ú
connectionr   r   r   rZ   ¶   r7   zTransport.close_connectionc                 C   r2   )NÚcreate_channelr5   r[   r   r   r   r]   ¹   r7   zTransport.create_channelc                 C   r2   )NÚclose_channelr5   r[   r   r   r   r^   ¼   r7   zTransport.close_channelc                 K   r2   )NÚdrain_eventsr5   )r-   r\   r/   r   r   r   r_   ¿   r7   zTransport.drain_eventsé   c                 C   ó   d S r   r   )r-   r\   Úrater   r   r   Úheartbeat_checkÂ   r<   zTransport.heartbeat_checkc                 C   r8   )NrW   r   r6   r   r   r   Údriver_versionÅ   r<   zTransport.driver_versionc                 C   r8   )Nr   r   r[   r   r   r   Úget_heartbeat_intervalÈ   r<   z Transport.get_heartbeat_intervalc                 C   ra   r   r   ©r-   r\   Úloopr   r   r   Úregister_with_event_loopË   r<   z"Transport.register_with_event_loopc                 C   ra   r   r   rf   r   r   r   Úunregister_from_event_loopÎ   r<   z$Transport.unregister_from_event_loopc                 C   r8   ©NTr   r[   r   r   r   Úverify_connectionÑ   r<   zTransport.verify_connectionc                    s    ˆj ‰‡ ‡‡‡‡‡fdd„‰ ˆ S )Nc              
      sr   ˆj stdƒ‚zˆdd� W n" ˆy   Y d S  ˆy0 } z|jˆv r+W Y d }~d S ‚ d }~ww |  ˆ | ¡ d S )NzSocket was disconnectedr   )Útimeout)Ú	connectedr   ÚerrnoÚ	call_soon)rg   Úexc©Ú_readÚ_unavailr\   r_   Úerrorrl   r   r   rr   Ø   s   
€ýz%Transport._make_reader.<locals>._read)r_   )r-   r\   rl   rt   rs   r   rq   r   Ú_make_readerÔ   s   zTransport._make_readerc                 C   r8   rj   r   r[   r   r   r   Úqos_semantics_matches_specç   r<   z$Transport.qos_semantics_matches_specc                 C   s*   | j }|d u r|  |¡ }| _ ||ƒ d S r   )Ú_Transport__readerru   )r-   r\   rg   Úreaderr   r   r   Úon_readableê   s   zTransport.on_readablec                 C   s   i S r   r   r6   r   r   r   Údefault_connection_paramsð   s   z#Transport.default_connection_paramsc                 O   s
   |   | ¡S r   )r   )r-   r.   r/   r   r   r   Úget_managerô   r7   zTransport.get_managerc                 C   s   |   ¡ S r   )r{   r6   r   r   r   Úmanager÷   ó   zTransport.managerc                 C   ó   | j jS r   )Ú
implementsrV   r6   r   r   r   Úsupports_heartbeatsû   r}   zTransport.supports_heartbeatsc                 C   r~   r   )r   rT   r6   r   r   r   Úsupports_evÿ   r}   zTransport.supports_ev)r`   ).rB   rC   rD   rE   r   rX   Úcan_parse_urlÚdefault_portr   Úconnection_errorsr   Úchannel_errorsÚdriver_typeÚdriver_namerw   Údefault_transport_capabilitiesrO   r   rI   rY   rZ   r]   r^   r_   rc   rd   re   rh   ri   rk   Úsocketrl   rt   rn   ÚEAGAINÚEINTRru   rv   ry   Úpropertyrz   r{   r
   r|   r€   r�   r   r   r   r   r   �   sL    

ÿ


r   )#rE   Ú
__future__r   r   rn   r‰   Úamqp.exceptionsr   Úkombu.exceptionsr   r   Ú
kombu.fiver   Úkombu.messager   Úkombu.utils.functionalr	   Úkombu.utils.objectsr
   Úkombu.utils.timer   Ú__all__Úintr!   r    r   r)   Úobjectr   r   r   rJ   Ú	frozensetrˆ   r   r   r   r   r   Ú<module>   s<    û	"#

ý