o
    —¨ÊhH  ã                   @  s°   d Z ddl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
 d	Zd
ZG dd„ de	jƒZG dd„ deje	jƒZG dd„ dejƒZG dd„ de	jƒZG dd„ deƒZdS )aÖ  pyamqp transport module for Kombu.

Pure-Python amqp transport using py-amqp library.

Features
========
* Type: Native
* Supports Direct: Yes
* Supports Topic: Yes
* Supports Fanout: Yes
* Supports Priority: Yes
* Supports TTL: Yes

Connection String
=================
Connection string can have the following formats:

.. code-block::

    amqp://[USER:PASSWORD@]BROKER_ADDRESS[:PORT][/VIRTUALHOST]
    [USER:PASSWORD@]BROKER_ADDRESS[:PORT][/VIRTUALHOST]
    amqp://

For TLS encryption use:

.. code-block::

    amqps://[USER:PASSWORD@]BROKER_ADDRESS[:PORT][/VIRTUALHOST]

Transport Options
=================
Transport Options are passed to constructor of underlying py-amqp
:class:`~kombu.connection.Connection` class.

Using TLS
=========
Transport over TLS can be enabled by ``ssl`` parameter of
:class:`~kombu.Connection` class. By setting ``ssl=True``, TLS transport is
used::

    conn = Connect('amqp://', ssl=True)

This is equivalent to ``amqps://`` transport URI::

    conn = Connect('amqps://')

For adding additional parameters to underlying TLS, ``ssl`` parameter should
be set with dict instead of True::

    conn = Connect('amqp://broker.example.com', ssl={
            'keyfile': '/path/to/keyfile'
            'certfile': '/path/to/certfile',
            'ca_certs': '/path/to/ca_certfile'
        }
    )

All parameters are passed to ``ssl`` parameter of
:class:`amqp.connection.Connection` class.

SSL option ``server_hostname`` can be set to ``None`` which is causing using
hostname from broker URL. This is usefull when failover is used to fill
``server_hostname`` with currently used broker::

    conn = Connect('amqp://broker1.example.com;broker2.example.com', ssl={
            'server_hostname': None
        }
    )
é    )ÚannotationsN)Úget_manager)Úversion_string_as_tupleé   )Úbase©Úto_rabbitmq_queue_argumentsi(  i'  c                      s"   e Zd ZdZd‡ fdd„	Z‡  ZS )ÚMessagezAMQP Message.Nc                   sL   |j }tƒ jd|j||j| d¡| d¡|j|j | d¡pi dœ|¤Ž d S )NÚcontent_typeÚcontent_encodingÚapplication_headers)ÚbodyÚchannelÚdelivery_tagr
   r   Údelivery_infoÚ
propertiesÚheaders© )r   ÚsuperÚ__init__r   r   Úgetr   )ÚselfÚmsgr   ÚkwargsÚprops©Ú	__class__r   úH/var/www/html/env/lib/python3.10/site-packages/kombu/transport/pyamqp.pyr   X   s   ø	
÷zMessage.__init__©N©Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   Ú__classcell__r   r   r   r   r	   U   s    r	   c                   @  s<   e Zd ZdZeZdddddejfdd„Zdd„ Zdd„ ZdS )	ÚChannelzAMQP Channel.Nc                 C  s   ||f||||dœ|pi ¤ŽS )z<Prepare message so that it can be sent using this transport.)Úpriorityr
   r   r   r   )r   r   r&   r
   r   r   r   Ú_Messager   r   r   Úprepare_messagek   s   ÿûúzChannel.prepare_messagec                 K  s   t |fi |¤ŽS r   r   )r   Ú	argumentsr   r   r   r   Úprepare_queue_argumentsx   ó   zChannel.prepare_queue_argumentsc                 C  s   | j || d�S )z4Convert encoded message body back to a Python value.©r   )r	   )r   Úraw_messager   r   r   Úmessage_to_python{   s   zChannel.message_to_python)	r    r!   r"   r#   r	   Úamqpr(   r*   r.   r   r   r   r   r%   f   s    
þr%   c                   @  s   e Zd ZdZeZdS )Ú
ConnectionzAMQP Connection.N)r    r!   r"   r#   r%   r   r   r   r   r0   €   s    r0   c                   @  sÐ   e Zd ZdZeZeZeZe	jj
Z
e	jjZe	jjZe	jjZdZdZejjjdd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„ Zd%dd„Zdd„ Ze d d!„ ƒZ!d"d#„ Z"dS )&Ú	TransportzAMQP Transport.zpy-amqpr/   T)ÚasynchronousÚ
heartbeatsNc                 K  s"   || _ |p| j| _|p| j| _d S r   )ÚclientÚdefault_portÚdefault_ssl_port)r   r4   r5   r6   r   r   r   r   r   ž   s   zTransport.__init__c                 C  s   t jS r   )r/   Ú__version__©r   r   r   r   Údriver_version¤   ó   zTransport.driver_versionc                 C  s   |  ¡ S r   r,   ©r   Ú
connectionr   r   r   Úcreate_channel§   s   zTransport.create_channelc                 K  s   |j di |¤ŽS )Nr   )Údrain_events)r   r<   r   r   r   r   r>   ª   r+   zTransport.drain_eventsc                 C  s   |d ur
|  ¡  d S d S r   )Úcollectr;   r   r   r   Ú_collect­   s   ÿzTransport._collectc                 C  sÒ   | j }| j ¡ D ]\}}t||dƒst|||ƒ q|jdkr!d|_t|jtƒr9d|jv r9|jd du r9|j|jd< t|j	|j
|j|j|j|j|j|j|jdœ	fi |jpTi ¤Ž}| jdi |¤Ž}| j |_ | ¡  |S )z(Establish connection to the AMQP broker.NÚ	localhostz	127.0.0.1Úserver_hostname)	ÚhostÚuseridÚpasswordÚlogin_methodÚvirtual_hostÚinsistÚsslÚconnect_timeoutÚ	heartbeatr   )r4   Údefault_connection_paramsÚitemsÚgetattrÚsetattrÚhostnameÚ
isinstancerI   ÚdictrC   rD   rE   rF   rG   rH   rJ   rK   Útransport_optionsr0   Úconnect)r   ÚconninfoÚnameÚdefault_valueÚoptsÚconnr   r   r   Úestablish_connection±   s8   €

÷
özTransport.establish_connectionc                 C  ó   |j S r   )Ú	connectedr;   r   r   r   Úverify_connectionÎ   r:   zTransport.verify_connectionc                 C  s   d|_ | ¡  dS )z!Close the AMQP broker connection.N)r4   Úcloser;   r   r   r   Úclose_connectionÑ   s   zTransport.close_connectionc                 C  r[   r   )rK   r;   r   r   r   Úget_heartbeat_intervalÖ   r:   z Transport.get_heartbeat_intervalc                 C  s    d|j _| |j| j||¡ d S ©NT)Ú	transportÚraise_on_initial_eintrÚ
add_readerÚsockÚon_readable)r   r<   Úloopr   r   r   Úregister_with_event_loopÙ   s   z"Transport.register_with_event_loopé   c                 C  s   |j |d�S )N)Úrate)Úheartbeat_tick)r   r<   rj   r   r   r   Úheartbeat_checkÝ   s   zTransport.heartbeat_checkc                 C  s(   |j }| d¡dkrt|d ƒdk S dS )NÚproductÚRabbitMQÚversion)é   rp   T)Úserver_propertiesr   r   )r   r<   r   r   r   r   Úqos_semantics_matches_specà   s   z$Transport.qos_semantics_matches_specc                 C  s    dd| j jr	| jn| jdddœS )NÚguestrA   ÚPLAIN)rD   rE   ÚportrP   rF   )r4   rI   r6   r5   r8   r   r   r   rL   æ   s   úz#Transport.default_connection_paramsc                 O  s   t | jg|¢R i |¤ŽS r   )r   r4   ©r   Úargsr   r   r   r   r   ñ   s   zTransport.get_manager)NN)ri   )#r    r!   r"   r#   r0   ÚDEFAULT_PORTr5   ÚDEFAULT_SSL_PORTr6   r/   Úconnection_errorsÚchannel_errorsÚrecoverable_connection_errorsÚrecoverable_channel_errorsÚdriver_nameÚdriver_typer   r1   Ú
implementsÚextendr   r9   r=   r>   r@   rZ   r]   r_   r`   rh   rl   rr   ÚpropertyrL   r   r   r   r   r   r1   †   s@    ÿþ
ÿ


r1   c                      s    e Zd ZdZ‡ fdd„Z‡  ZS )ÚSSLTransportzAMQP SSL Transport.c                   s*   t ƒ j|i |¤Ž | jjsd| j_d S d S ra   )r   r   r4   rI   rv   r   r   r   r   ø   s   ÿzSSLTransport.__init__r   r   r   r   r   rƒ   õ   s    rƒ   )r#   Ú
__future__r   r/   Úkombu.utils.amq_managerr   Úkombu.utils.textr   Ú r   r   rx   ry   r	   r%   Ú
StdChannelr0   r1   rƒ   r   r   r   r   Ú<module>   s    Fo