o
    `j                     @   s   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ejejhZdZG dd dZG dd	 d	e
ZdS )
    N)urlparse)KafkaConnectionError)KafkaNetSocketi   c                   @   s    e Zd ZdZdZdZdZdZdS )_Statesz<disconnected>z<connecting>z	<sending>z	<reading>z
<complete>N)__name__
__module____qualname__DISCONNECTED
CONNECTINGSENDINGREADINGCOMPLETE r   r   T/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/net/http_connect.pyr      s    r   c                       sl   e Zd ZdZ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 )HttpConnectProxya}  Tunnels broker connections through an HTTP CONNECT proxy (RFC 7231 s4.3.6).

    Registered for the ``http`` scheme -- pass ``proxy_url='http://host:port'``
    to KafkaConsumer/KafkaProducer/KafkaAdminClient.

    Basic proxy auth is supported via URL credentials: ``http://user:pass@host:8080``.
    Broker hostnames are always forwarded unresolved so the proxy handles DNS.
    )httpc                 C   s2   t || _d | _tj| _d| _d| _|  | _	d S )N    )
r   
_proxy_url_sockr   r	   _state	_send_buf	_recv_buf_get_proxy_addr_proxy_addr)self	proxy_urlr   r   r   __init__&   s   
zHttpConnectProxy.__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   addrsr   r   r   r   .   s   
z HttpConnectProxy._get_proxy_addrFc                    s0   |rt  j||ddS tjtjtjd||ffgS )NT)raise_error )superr   socket	AF_UNSPECSOCK_STREAMIPPROTO_TCP)r   hostr    r   	__class__r   r   r   4   s   zHttpConnectProxy.dns_lookupc                 C   s,   || _ | j\}}}}}t|||| _| jS )N)_target_afir   r'   r   )r   family	sock_typeprotoproxy_family_r   r   r   r'   :   s   zHttpConnectProxy.socketc                 C   s   || j u sJ | jtjkrtj| _| jtjkr"| |}|d ur"|S | jtjkr2|  }|d ur2|S | jtjkrB| 	 }|d urB|S | jtj
krJdS tjS )Nr   )r   r   r   r	   r
   _do_connectingr   _do_sendingr   _do_readingr   errnoECONNREFUSED)r   sockaddrretr   r   r   
connect_ex@   s$   
zHttpConnectProxy.connect_exc           	      C   s   | j \}}}}}| j|}|r|tjkr|S |d |d }}d||}| jjrF| jjrFt	
d| jj| jj  }|d|7 }|d  | _tj| _d S )Nr      z)CONNECT {0}:{1} HTTP/1.1
Host: {0}:{1}
z{0}:{1}zProxy-Authorization: Basic {}
z
)r   r   r<   r7   EISCONNformatr   usernamepasswordbase64	b64encodeencodedecoder   r   r   r   )	r   r:   r3   proxy_sockaddrr;   r+   r    headerscredentialsr   r   r   r4   Z   s    zHttpConnectProxy._do_connectingc              
   C   s   | j r@z| j| j }|dkrtd tjW S | j |d  | _ W n ty< } z|jtv r7tj	W  Y d }~S  d }~ww | j st
j| _d S )Nr   z5Proxy closed connection while sending CONNECT request)r   r   sendlogerrorr7   r8   OSError_WOULD_BLOCKEWOULDBLOCKr   r   r   )r   sentexcr   r   r   r5   j   s    

zHttpConnectProxy._do_sendingc              
   C   s   d| j vr[z5| jd}|std | j  tjW S |  j |7  _ t| j t	kr9tdt	 | j  tjW S W n t
yU } z|jtv rPtjW  Y d }~S  d }~ww d| j vs| j dd }d|v sl|drrtj| _d S td	| | j  tjS )
Ns   

i   z0Proxy closed connection during CONNECT handshakez7Proxy response exceeded %d bytes without end-of-headerss   
r   s    200 s    200z HTTP CONNECT to proxy failed: %r)r   r   recvrJ   rK   closer7   r8   len_MAX_RESPONSE_SIZErL   rM   rN   splitendswithr   r   r   )r   chunkrP   
first_liner   r   r   r6   y   s6   






zHttpConnectProxy._do_reading)F)r   r   r   __doc__SCHEMESr   r   r   r'   r(   r)   r*   r<   r4   r5   r6   __classcell__r   r   r,   r   r      s    	r   )rB   r7   loggingr!   r'   urllib.parser   kafka.errorsr   kafka.net.inetr   	getLoggerr   rJ   rN   EAGAINrM   rT   r   r   r   r   r   r   <module>   s    
