o
    `šjy(  ã                   @   sx   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
 e e¡ZG dd„ dƒZG dd„ de
ƒZdS )	é    N)Úurlparse)ÚKafkaConnectionError)ÚKafkaNetSocketc                   @   s0   e Zd ZdZdZdZdZdZdZdZ	dZ
d	Zd
S )ÚProxyConnectionStatesz<disconnected>z<connecting>z<negotiate_propose>z<negotiating>z<authenticating>z<request_submit>z<requesting>z<read_address>z
<complete>N)Ú__name__Ú
__module__Ú__qualname__ÚDISCONNECTEDÚ
CONNECTINGÚNEGOTIATE_PROPOSEÚNEGOTIATINGÚAUTHENTICATINGÚREQUEST_SUBMITÚ
REQUESTINGÚREAD_ADDRESSÚCOMPLETE© r   r   úN/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/net/socks5.pyr      s    r   c                       st   e Zd ZdZdZdd„ Zdd„ Zdd„ Zd‡ fd
d„	Ze	j
e	je	jfdd„Z	dd„ Zdd„ Zdd„ Zdd„ Z‡  ZS )ÚSocks5ProxyzuSocks5 proxy

    Manages connection through socks5 proxy with support for username/password
    authentication.
    )Úsocks5Úsocks5hc                 C   sZ   d| _ d| _t|ƒ| _| jj| jvrtd| jjf ƒ‚d | _tj	| _
tj| _|  ¡ | _d S )Nó    zUnsupported proxy scheme: %s)Ú
_buffer_inÚ_buffer_outr   Ú
_proxy_urlÚschemeÚSCHEMESÚ
ValueErrorÚ_sockr   r	   Ú_stateÚsocketÚ	AF_UNSPECÚ_target_afiÚ_get_proxy_addrÚ_proxy_addr)ÚselfÚ	proxy_urlr   r   r   Ú__init__$   s   
zSocks5Proxy.__init__c                 C   s.   | j | jj| jjdd}|stdƒ‚t |¡S )NT)Úproxyz#Unable to resolve proxy_url via dns)Ú
dns_lookupr   ÚhostnameÚportr   ÚrandomÚchoice)r%   Úproxy_addrsr   r   r   r#   /   s   
zSocks5Proxy._get_proxy_addrc                 C   s   | j jdkS )Nr   )r   r   )r%   r   r   r   Ú_use_remote_lookup5   s   zSocks5Proxy._use_remote_lookupFc                    sF   |rt ƒ j||ddS |  ¡ rtjtjtjd||ffgS t ƒ  ||¡S )NT)Úraise_errorÚ )Úsuperr)   r/   r    r!   ÚSOCK_STREAMÚIPPROTO_TCP)r%   Úhostr+   r(   ©Ú	__class__r   r   r)   8   s
   zSocks5Proxy.dns_lookupc                 C   s,   || _ | j\}}}}}t |||¡| _| jS )zšOpen and record a socket.

        Returns the actual underlying socket
        object to ensure e.g. selects and ssl wrapping works as expected.
        )r"   r$   r    r   )r%   ÚfamilyÚ	sock_typeÚprotoÚproxy_familyÚ_r   r   r   r    @   s   zSocks5Proxy.socketc                 C   s2   | j r| j | j ¡}| j |d… | _ | j sdS dS )zÊSend out all data that is stored in the outgoing buffer.

        It is expected that the caller handles error handling, including non-blocking
        as well as connection failure exceptions.
        N)r   r   Úsend)r%   Ú
sent_bytesr   r   r   Ú
_flush_bufK   s   þzSocks5Proxy._flush_bufc                 C   sH   	 |t | jƒ }|dkrn| j |¡}|sn| j| | _q| jd|… S )z´Ensure local inbound buffer has enough data, and return that data without
        consuming the local buffer

        It's expected that the caller handles e.g. blocking exceptionsTr   N)Úlenr   r   Úrecv)r%   ÚdatalenÚbytes_remainingÚdatar   r   r   Ú	_peek_bufU   s   ù	zSocks5Proxy._peek_bufc                 C   s&   |   |¡}|r| jt|ƒd… | _|S )zuRead and consume bytes from socket connection

        It's expected that the caller handles e.g. blocking exceptionsN)rE   r   r@   )r%   rB   Úbufr   r   r   Ú	_read_bufe   s   
zSocks5Proxy._read_bufc              
   C   s  || j u sJ ‚| jtjkrtj| _| jtjkr3| j\}}}}}| j  |¡}|r,|tjkr1tj	| _n|S | jtj	krL| j
jrE| j
jrEd| _nd| _tj| _| jtjkrÂ|  ¡  |  d¡}|dd… dkrtt d¡ tj| _| j  ¡  tjS |dd… dkrtj| _nA|dd… d	kr±t| j
jƒ}t| j
jƒ}t d
 ||¡d|| j
j ¡ || j
j ¡ ¡| _tj| _nt d¡ tj| _| j  ¡  tjS | jtjkrë|  ¡  |  d¡}|dkrÚtj| _nt d¡ tj| _| j  ¡  tjS | jtjkrt|  ¡ rÿd}	t|d ƒ}
n+| jtjkrd}	d}
n| jtj krd}	d}
nt d| j¡ tj| _| j  ¡  tjS t dddd|	¡| _|	dkrN|  jt d |
¡|
|d  d¡¡7  _n|  jt d |
¡t !| j|d ¡¡7  _|  jt d|d ¡7  _tj"| _| jtj"kr¨|  ¡  |  d¡}|dd… dkr’tj#| _nt d|dd… ¡ tj| _| j  ¡  tjS | jtj#krì|  $d¡}|dd… dkrÃ|  d¡}n%|dd… dkrÒ|  d¡}nt d|dd… ¡ tj| _| j  ¡  tjS tj%| _| jtj%krõdS t d| j¡ tj| _| j r	| j  ¡  tjS ) a
  Runs a state machine through connection to authentication to
        proxy connection request.

        The somewhat strange setup is to facilitate non-intrusive use from
        BrokerConnection state machine.

        This function is called with a socket in non-blocking mode. Both
        send and receive calls can return in EWOULDBLOCK/EAGAIN which we
        specifically avoid handling here. These are handled in main
        BrokerConnection connection loop, which then would retry calls
        to this function.s   s    é   r   é   ó   zUnrecognized SOCKS versionó    ó   z
!bb{}sb{}sz(Unrecognized SOCKS authentication methods    z#Socks5 proxy authentication failureé   é   é   zUnknown address family, %rz!bbbbé   z!b{}sÚasciiz!{}sz!Hs    zProxy request failed: %rs    é   s    é   z#Unrecognized remote address type %rz.Internal error, state %r not handled correctly)&r   r   r   r	   r
   r$   Ú
connect_exÚerrnoÚEISCONNr   r   ÚusernameÚpasswordr   r   r?   rG   ÚlogÚerrorÚcloseÚECONNREFUSEDr   r@   ÚstructÚpackÚformatÚencoder   r/   r"   r    ÚAF_INETÚAF_INET6Ú	inet_ptonr   r   rE   r   )r%   ÚsockÚaddrr<   ÚsockaddrÚretrF   ÚuserlenÚpasslenÚ	addr_typeÚaddr_lenr   r   r   rT   n   sÚ   







ú







û


ý
þ





zSocks5Proxy.connect_ex)F)r   r   r   Ú__doc__r   r'   r#   r/   r)   r    r!   r3   r4   r?   rE   rG   rT   Ú__classcell__r   r   r6   r   r      s    
	r   )rU   Úloggingr,   r    r]   Úurllib.parser   Úkafka.errorsr   Úkafka.net.inetr   Ú	getLoggerr   rY   r   r   r   r   r   r   Ú<module>   s    
