o
    Œ¨Êh¥Y  ã                   @   s  d Z ddl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 ddl	m
Z
mZ ddlmZ ddlmZmZ dd	lmZ ejejejejhZd
Zeƒ ZdZdZe d¡ZddddddœZefdd„Z G dd„ dƒZ!G dd„ de!ƒZ"G dd„ de!ƒZ#ddd„Z$dS )zTransport implementation.é    N)Úcontextmanager)ÚSSLError)ÚpackÚunpacké   )ÚUnexpectedFrame)ÚKNOWN_TCP_OPTSÚSOL_TCP)Úset_cloexeci(  iÿÿÿs   AMQP  	z\[([\.0-9a-f:]+)\](?::(\d+))?iè  é<   é
   é	   )ÚTCP_NODELAYÚTCP_USER_TIMEOUTÚTCP_KEEPIDLEÚTCP_KEEPINTVLÚTCP_KEEPCNTc                 C   sd   |}t  | ¡}|r| d¡} | d¡rt| d¡ƒ}| |fS d| v r.|  dd¡\} }t|ƒ}| |fS )z1Convert hostname:port string to host, port tuple.r   é   ú:)ÚIPV6_LITERALÚmatchÚgroupÚintÚrsplit)ÚhostÚdefaultÚportÚm© r   ú@/var/www/html/env/lib/python3.10/site-packages/amqp/transport.pyÚto_host_port(   s   


ýr    c                   @   sž   e Zd ZdZ			d$dd„ZdZdd„ Zd	d
„ Ze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fd d!„Zd"d#„ ZdS )&Ú_AbstractTransportaÄ  Common superclass for TCP and SSL transports.

    PARAMETERS:
        host: str

            Broker address in format ``HOSTNAME:PORT``.

        connect_timeout: int

            Timeout of creating new connection.

        read_timeout: int

            sets ``SO_RCVTIMEO`` parameter of socket.

        write_timeout: int

            sets ``SO_SNDTIMEO`` parameter of socket.

        socket_settings: dict

            dictionary containing `optname` and ``optval`` passed to
            ``setsockopt(2)``.

        raise_on_initial_eintr: bool

            when True, ``socket.timeout`` is raised
            when exception is received during first read. See ``_read()`` for
            details.
    NTc                 K   sD   d| _ d | _|| _t| _t|ƒ\| _| _|| _|| _	|| _
|| _d S ©NF)Ú	connectedÚsockÚraise_on_initial_eintrÚEMPTY_BUFFERÚ_read_bufferr    r   r   Úconnect_timeoutÚread_timeoutÚwrite_timeoutÚsocket_settings)Úselfr   r(   r)   r*   r+   r%   Úkwargsr   r   r   Ú__init__W   s   
z_AbstractTransport.__init__)Ú
connectionr$   r%   r'   r   r   r(   r)   r*   r+   Ú__dict__Ú__weakref__c              	   C   s’   | j r:| j  ¡ d › d| j  ¡ d › �}| j  ¡ d › d| j  ¡ d › �}dt| ƒj› d|› d|› dt| ƒd›d	�	S dt| ƒj› d
t| ƒd›d	�S )Nr   r   r   ú<z: z -> z at z#xú>z: (disconnected) at )r$   ÚgetsocknameÚgetpeernameÚtypeÚ__name__Úid)r,   ÚsrcÚdstr   r   r   Ú__repr__t   s
   ""*z_AbstractTransport.__repr__c              	   C   sr   z | j rW d S |  | j| j| j¡ |  | j| j| j¡ d| _ W d S  t	t
fy8   | jr7| j s7| j ¡  d | _‚ w )NT)r#   Ú_connectr   r   r(   Ú_init_socketr+   r)   r*   ÚOSErrorr   r$   Úclose©r,   r   r   r   Úconnect|   s   ÿ
ûz_AbstractTransport.connectc              
   c   sæ   � |d u r| j V  d S | j }| ¡ }||kr| |¡ zLz| j V  W n7 tyC } zdt|ƒv r4t ¡ ‚dt|ƒv r>t ¡ ‚‚ d }~w tyY } z|jtj	krTt ¡ ‚‚ d }~ww W ||krf| |¡ d S d S ||krr| |¡ w w )Nú	timed outzThe operation did not complete)
r$   Ú
gettimeoutÚ
settimeoutr   ÚstrÚsocketÚtimeoutr>   ÚerrnoÚEWOULDBLOCK)r,   rG   r$   ÚprevÚexcr   r   r   Úhaving_timeout�   s8   €
€€ý÷ÿÿz!_AbstractTransport.having_timeoutc              	   C   sÊ   t  ||t jt jt¡}t|ƒD ]S\}}|\}}}	}
}z*t   |||	¡| _zt| jdƒ W n	 ty4   Y nw | j 	|¡ | j 
|¡ W  d S  t jyb   | jrT| j ¡  d | _|d t|ƒkr`‚ Y qw d S )NTr   )rF   ÚgetaddrinfoÚ	AF_UNSPECÚSOCK_STREAMr	   Ú	enumerater$   r
   ÚNotImplementedErrorrD   rA   Úerrorr?   Úlen)r,   r   r   rG   ÚentriesÚiÚresÚafÚsocktypeÚprotoÚ	canonnameÚsar   r   r   r<   «   s0   ÿÿù
ÿüöz_AbstractTransport._connectc              	   C   s˜   | j  d ¡ | j  tjtjd¡ |  |¡ tj|ftj|ffD ]!\}}|d ur@t	|ƒ}t	|| d ƒ}| j  tj|t
d||ƒ¡ q|  ¡  |  t¡ d S )Nr   i@B Úll)r$   rD   Ú
setsockoptrF   Ú
SOL_SOCKETÚSO_KEEPALIVEÚ_set_socket_optionsÚSO_SNDTIMEOÚSO_RCVTIMEOr   r   Ú_setup_transportÚ_writeÚAMQP_PROTOCOL_HEADER)r,   r+   r)   r*   rG   ÚintervalÚsecÚusecr   r   r   r=   Â   s    
ÿ
þ€z_AbstractTransport._init_socketc              	   C   s”   i }t D ]C}d }|dkr zddlm} W n ty   d}Y nw tt|ƒr*tt|ƒ}|rG|tv r7t| ||< qtt|ƒrG| ttt|ƒ¡||< q|S )Nr   r   )r   é   )	r   rF   r   ÚImportErrorÚhasattrÚgetattrÚDEFAULT_SOCKET_SETTINGSÚ
getsockoptr	   )r,   r$   Útcp_optsÚoptÚenumr   r   r   Ú_get_tcp_socket_defaultsÕ   s(   þ



ÿ€z+_AbstractTransport._get_tcp_socket_defaultsc                 C   s@   |   | j¡}|r| |¡ | ¡ D ]\}}| j t||¡ qd S ©N)rr   r$   ÚupdateÚitemsr]   r	   )r,   r+   ro   rp   Úvalr   r   r   r`   ê   s   
ÿz&_AbstractTransport._set_socket_optionsFc                 C   ó   t dƒ‚)z#Read exactly n bytes from the peer.úMust be overridden in subclass©rQ   )r,   ÚnÚinitialr   r   r   Ú_readñ   ó   z_AbstractTransport._readc                 C   ó   dS )z.Do any additional initialization of the class.Nr   r@   r   r   r   rc   õ   ó   z#_AbstractTransport._setup_transportc                 C   r~   )z8Do any preliminary work in shutting down the connection.Nr   r@   r   r   r   Ú_shutdown_transportù   r   z&_AbstractTransport._shutdown_transportc                 C   rw   )z&Completely write a string to the peer.rx   ry   )r,   Úsr   r   r   rd   ý   r}   z_AbstractTransport._writec                 C   s‚   | j d ur<z|  ¡  W n	 ty   Y nw z	| j  tj¡ W n	 ty'   Y nw z| j  ¡  W n	 ty8   Y nw d | _ d| _d S r"   )r$   r€   r>   ÚshutdownrF   Ú	SHUT_RDWRr?   r#   r@   r   r   r   r?     s$   
ÿÿÿ
z_AbstractTransport.closec              
   C   sn  | j }t}zJ|ddƒ}||7 }|d|ƒ\}}}|tkr@|tƒ}z||t ƒ}	W n tjttfy7   ||7 }‚ w d ||	g¡}
n||ƒ}
||
7 }t|dƒƒ}W nU tjy^   || j	 | _	‚  ttfy¤ } z9t
|tjƒr‚tjdkr‚|jtjkr‚|| j	 | _	t ¡ ‚t
|tƒr—dt|ƒv r—|| j	 | _	t ¡ ‚|jtvrŸd| _‚ d	}~ww |d
kr®|||
fS td|d›d�ƒ‚)a¸  Parse AMQP frame.

        Frame has following format::

            0      1         3         7                   size+7      size+8
            +------+---------+---------+   +-------------+   +-----------+
            | type | channel |  size   |   |   payload   |   | frame-end |
            +------+---------+---------+   +-------------+   +-----------+
             octet    short     long        'size' octets        octet

        é   Tz>BHIó    r   ÚntrB   FNéÎ   zReceived frame_end z#04xz while expecting 0xce)r|   r&   ÚSIGNED_INT_MAXrF   rG   r>   r   ÚjoinÚordr'   Ú
isinstancerR   ÚosÚnamerH   rI   rE   Ú_UNAVAILr#   r   )r,   r   ÚreadÚread_frame_bufferÚframe_headerÚ
frame_typeÚchannelÚsizeÚpart1Úpart2ÚpayloadÚ	frame_endrK   r   r   r   Ú
read_frame  sR   
ü
ÿ

€í
ÿz_AbstractTransport.read_framec              
   C   sL   z|   |¡ W d S  tjy   ‚  ty% } z	|jtvr d| _‚ d }~ww r"   )rd   rF   rG   r>   rH   rŽ   r#   )r,   r�   rK   r   r   r   ÚwriteY  s   
€ýz_AbstractTransport.write)NNNNT)F)r7   Ú
__module__Ú__qualname__Ú__doc__r.   Ú	__slots__r;   rA   r   rL   r<   r=   rr   r`   r|   rc   r€   rd   r?   r   r™   rš   r   r   r   r   r!   7   s,    
þ

Br!   c                       s€   e Zd ZdZd‡ fdd„	ZdZdd„ Zddd	„Zdd
d„Z					ddd„Z	dd„ Z
dejejejffdd„Zdd„ Z‡  ZS )ÚSSLTransportaÏ  Transport that works over SSL.

    PARAMETERS:
        host: str

            Broker address in format ``HOSTNAME:PORT``.

        connect_timeout: int

            Timeout of creating new connection.

        ssl: bool|dict

            parameters of TLS subsystem.
                - when ``ssl`` is not dictionary, defaults of TLS are used
                - otherwise:
                    - if ``ssl`` dictionary contains ``context`` key,
                      :attr:`~SSLTransport._wrap_context` is used for wrapping
                      socket. ``context`` is a dictionary passed to
                      :attr:`~SSLTransport._wrap_context` as context parameter.
                      All others items from ``ssl`` argument are passed as
                      ``sslopts``.
                    - if ``ssl`` dictionary does not contain ``context`` key,
                      :attr:`~SSLTransport._wrap_socket_sni` is used for
                      wrapping socket. All items in ``ssl`` argument are
                      passed to :attr:`~SSLTransport._wrap_socket_sni` as
                      parameters.

        kwargs:

            additional arguments of
            :class:`~amqp.transport._AbstractTransport` class
    Nc                    s6   t |tƒr|ni | _t| _tƒ j|fd|i|¤Ž d S )Nr(   )r‹   ÚdictÚssloptsr&   r'   Úsuperr.   )r,   r   r(   Ússlr-   ©Ú	__class__r   r   r.   ‡  s   ÿÿ
ÿzSSLTransport.__init__)r¡   c                 C   s>   | j | jfi | j¤Ž| _| j | j¡ | j ¡  | jj| _dS )z!Wrap the socket in an SSL object.N)Ú_wrap_socketr$   r¡   rD   r(   Údo_handshaker�   Ú_quick_recvr@   r   r   r   rc   ‘  s   
zSSLTransport._setup_transportc                 K   s*   |r| j ||fi |¤ŽS | j|fi |¤ŽS rs   )Ú_wrap_contextÚ_wrap_socket_sni)r,   r$   Úcontextr¡   r   r   r   r¦   ™  s   zSSLTransport._wrap_socketc                 K   s(   t jdi |¤Ž}||_|j|fi |¤ŽS )uÝ  Wrap socket without SNI headers.

        PARAMETERS:
            sock: socket.socket

            Socket to be wrapped.

            sslopts: dict

                Parameters of  :attr:`ssl.SSLContext.wrap_socket`.

            check_hostname

                Whether to match the peer certâ€™s hostname. See
                :attr:`ssl.SSLContext.check_hostname` for details.

            ctx_options

                Parameters of :attr:`ssl.create_default_context`.
        Nr   )r£   Úcreate_default_contextÚcheck_hostnameÚwrap_socket)r,   r$   r¡   r­   Úctx_optionsÚctxr   r   r   r©   ž  s   zSSLTransport._wrap_contextFTc                 C   sæ   |||||	dœ}|du r|rt jnt j}t  |¡}|dur#| ||¡ |dur,| |¡ |
dur5| |
¡ z
t jo<|	du|_W n	 t	yH   Y nw |durP||_
|du ri|j
t jkri|r`t jjnt jj}| |¡ |jdi |¤Ž}|S )u�  Socket wrap with SNI headers.

        stdlib :attr:`ssl.SSLContext.wrap_socket` method augmented with support
        for setting the server_hostname field required for SNI hostname header.

        PARAMETERS:
            sock: socket.socket

                Socket to be wrapped.

            keyfile: str

                Path to the private key

            certfile: str

                Path to the certificate

            server_side: bool

                Identifies whether server-side or client-side
                behavior is desired from this socket. See
                :attr:`~ssl.SSLContext.wrap_socket` for details.

            cert_reqs: ssl.VerifyMode

                When set to other than :attr:`ssl.CERT_NONE`, peers certificate
                is checked. Possible values are :attr:`ssl.CERT_NONE`,
                :attr:`ssl.CERT_OPTIONAL` and :attr:`ssl.CERT_REQUIRED`.

            ca_certs: str

                Path to â€œcertification authorityâ€� (CA) certificates
                used to validate other peersâ€™ certificates when ``cert_reqs``
                is other than :attr:`ssl.CERT_NONE`.

            do_handshake_on_connect: bool

                Specifies whether to do the SSL
                handshake automatically. See
                :attr:`~ssl.SSLContext.wrap_socket` for details.

            suppress_ragged_eofs (bool):

                See :attr:`~ssl.SSLContext.wrap_socket` for details.

            server_hostname: str

                Specifies the hostname of the service which
                we are connecting to. See :attr:`~ssl.SSLContext.wrap_socket`
                for details.

            ciphers: str

                Available ciphers for sockets created with this
                context. See :attr:`ssl.SSLContext.set_ciphers`

            ssl_version:

                Protocol of the SSL Context. The value is one of
                ``ssl.PROTOCOL_*`` constants.
        )r$   Úserver_sideÚdo_handshake_on_connectÚsuppress_ragged_eofsÚserver_hostnameNr   )r£   ÚPROTOCOL_TLS_SERVERÚPROTOCOL_TLS_CLIENTÚ
SSLContextÚload_cert_chainÚload_verify_locationsÚset_ciphersÚHAS_SNIr­   ÚAttributeErrorÚverify_modeÚ	CERT_NONEÚPurposeÚCLIENT_AUTHÚSERVER_AUTHÚload_default_certsr®   )r,   r$   ÚkeyfileÚcertfiler±   Ú	cert_reqsÚca_certsr²   r³   r´   ÚciphersÚssl_versionÚoptsr«   Úpurposer   r   r   rª   ·  sD   Dûÿý


ÿÿ
ÿý
zSSLTransport._wrap_socket_snic                 C   s   | j dur| j  ¡ | _ dS dS )z/Unwrap a SSL socket, so we can call shutdown().N)r$   Úunwrapr@   r   r   r   r€   /  s   
ÿz SSLTransport._shutdown_transportc           	   
   C   sÄ   | j }| j}zDt|ƒ|k rIz
||t|ƒ ƒ}W n! ty8 } z|j|v r3|r-| jr-t ¡ ‚W Y d }~q‚ d }~ww |s?tdƒ‚||7 }t|ƒ|k sW n   || _‚ |d |… ||d … }| _|S )Nú%Server unexpectedly closed connection©r¨   r'   rS   r>   rH   r%   rF   rG   ©	r,   rz   r{   Ú_errnosÚrecvÚrbufr�   rK   Úresultr   r   r   r|   4  s0   

€ùó€zSSLTransport._readc                 C   sT   | j j}|r(z||ƒ}W n ty   d}Y nw |stdƒ‚||d… }|sdS dS )z+Write a string out to the SSL socket fully.r   zSocket closedN)r$   rš   Ú
ValueErrorr>   )r,   r�   rš   rz   r   r   r   rd   P  s   ûõzSSLTransport._write)NNrs   )
NNFNNFTNNN)r7   r›   rœ   r�   r.   rž   rc   r¦   r©   rª   r€   rH   ÚENOENTÚEAGAINÚEINTRr|   rd   Ú__classcell__r   r   r¤   r   rŸ   d  s$    "


üx
ÿrŸ   c                   @   s.   e Zd ZdZdd„ Zdejejffdd„ZdS )ÚTCPTransportz~Transport that deals directly with TCP socket.

    All parameters are :class:`~amqp.transport._AbstractTransport` class.
    c                 C   s   | j j| _t| _| j j| _d S rs   )r$   Úsendallrd   r&   r'   rÐ   r¨   r@   r   r   r   rc   g  s   
zTCPTransport._setup_transportFc           	   
   C   sÄ   | j }| j}zDt|ƒ|k rIz
||t|ƒ ƒ}W n! ty8 } z|j|v r3|r-| jr-t ¡ ‚W Y d}~q‚ d}~ww |s?tdƒ‚||7 }t|ƒ|k sW n   || _‚ |d|… ||d… }| _|S )z%Read exactly n bytes from the socket.NrÌ   rÍ   rÎ   r   r   r   r|   n  s0   

€ûõ€zTCPTransport._readN)	r7   r›   rœ   r�   rc   rH   rÕ   rÖ   r|   r   r   r   r   rØ   a  s    rØ   Fc                 K   s"   |rt nt}|| f||dœ|¤ŽS )a›  Create transport.

    Given a few parameters from the Connection constructor,
    select and create a subclass of
    :class:`~amqp.transport._AbstractTransport`.

    PARAMETERS:

        host: str

            Broker address in format ``HOSTNAME:PORT``.

        connect_timeout: int

            Timeout of creating new connection.

        ssl: bool|dict

            If set, :class:`~amqp.transport.SSLTransport` is used
            and ``ssl`` parameter is passed to it. Otherwise
            :class:`~amqp.transport.TCPTransport` is used.

        kwargs:

            additional arguments of :class:`~amqp.transport._AbstractTransport`
            class
    )r(   r£   )rŸ   rØ   )r   r(   r£   r-   Ú	transportr   r   r   Ú	Transport‡  s   rÛ   r"   )%r�   rH   rŒ   ÚrerF   r£   Ú
contextlibr   r   Ústructr   r   Ú
exceptionsr   Úplatformr   r	   Úutilsr
   rÕ   rÖ   rÔ   rI   rŽ   Ú	AMQP_PORTÚbytesr&   rˆ   re   Úcompiler   rm   r    r!   rŸ   rØ   rÛ   r   r   r   r   Ú<module>   s@    
û	  / ~&