o
    `j                     @   s`   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dS )    N)KafkaConnectionManager)NetworkSelectorc                   @   s   e Zd ZdZd2ddZedd Zdd Zd	d
 Zdd Z	dd Z
dd Zd3ddZd4ddZdd Zdd Zd4ddZd5ddZdd  Zd3d!d"Zd2d#d$Zd4d%d&Zd6d(d)Zd6d*d+Zd,d- Zd.d/ Zd4d0d1ZdS )7KafkaNetClienta3  Drop-in replacement for KafkaClient backed by KafkaConnectionManager.

    Provides the KafkaClient API surface that existing consumer/producer/admin
    code depends on. Goal: shrink over time as components transition to using
    KafkaConnectionManager directly (fire-and-forget via _request_buffer).
    Nc                 K   sP   t  | _|d u rtdi |n|| _|d u r#t| jfi || _d S || _d S )N )	threadingRLock_lockr   _netr   _manager)selfnetmanagerconfigsr   r   N/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/net/compat.py__init__   s   
*zKafkaNetClient.__init__c                 C   s   | j jS N)r
   clusterr   r   r   r   r      s   zKafkaNetClient.clusterc                 C   s   | j j|}|d uo|jS r   )r
   _connsget	connectedr   node_idconnr   r   r   r   "   s   zKafkaNetClient.connectedc                 C   s   |  | S r   )r   r   r   r   r   r   is_disconnected&      zKafkaNetClient.is_disconnectedc                 C   s$   | j j|}|d uo|jo|j S r   )r
   r   r   r   pausedr   r   r   r   is_ready)   s   zKafkaNetClient.is_readyc                 K   s8   |  |rdS z	| j| W dS  tjy   Y dS w )NTF)r   r
   get_connectionErrorsNodeNotReadyErrorr   r   kwargsr   r   r   ready-   s   
zKafkaNetClient.readyc                 K   s*   z	| j | W d S  tjy   Y d S w r   )r
   r   r    r!   r"   r   r   r   maybe_connect6   s
   zKafkaNetClient.maybe_connect0u  c                 C   s   |  |rdS | | | jj|}|d ur$|jjs$| jj||jd |d ur<|j	r<|j
r<| jjt|| jjd d |  |sJtd||f dS )NT
timeout_msfuturerequest_timeout_msr(   zNode %s not ready after %s ms)r   r%   r
   r   r   init_futureis_doner	   pollr   r   minconfigr    KafkaConnectionError)r   r   r(   r   r   r   r   await_ready<   s   


zKafkaNetClient.await_readyc                 C   sF   |d ur| j j|}|d urt|jS dS tdd | j j D S )Nr   c                 s   s    | ]}t |jV  qd S r   )lenin_flight_requests).0cr   r   r   	<genexpr>Q   s    z9KafkaNetClient.in_flight_request_count.<locals>.<genexpr>)r
   r   r   r3   r4   sumvaluesr   r   r   r   in_flight_request_countM   s   z&KafkaNetClient.in_flight_request_countc                 C   s6   | j j|}|d u rdS |jt  }td|d S )Nr     )r
   r   r   _throttle_timetime	monotonicmax)r   r   r   	remainingr   r   r   throttle_delayS   s
   zKafkaNetClient.throttle_delayc                 C   s   | j j}|d uo|j S r   )r
   _bootstrap_futurer-   )r   bootstrap_futurer   r   r   bootstrap_connected\   s   z"KafkaNetClient.bootstrap_connectedc                 C   s    | j jd u r| j | | j jS r   )r
   broker_version	bootstrap)r   r(   r   r   r   get_broker_version`   s   z!KafkaNetClient.get_broker_version'  c                    s@    j js
 j | |d u r j jS  fdd} j|||S )Nc                    s    j j| |dI d H }|jS )Nr+   )r
   r   rE   )	broker_idr(   r   r   r   r   _check_versionj   s   z4KafkaNetClient.check_version.<locals>._check_version)r
   bootstrappedrF   rE   r	   run)r   r   r(   rJ   r   r   r   check_versione   s   zKafkaNetClient.check_versionc                 K   s   | j j||dS N)r   )r
   send)r   r   requestr#   r   r   r   rO   q   s   zKafkaNetClient.sendc                 C   sP   | j ||d | ||}| jj||d | r|jS | r#|jt	d)Nr+   r'   zRequest timed out)
r2   rO   r	   r.   	succeededvaluefailed	exceptionr    KafkaTimeoutError)r   r   rP   r(   fr   r   r   send_and_receivet   s   
zKafkaNetClient.send_and_receivec                 C   s:   | j  | jj||dW  d    S 1 sw   Y  d S )Nr'   )r   r	   r.   )r   r(   r)   r   r   r   r.      s   $zKafkaNetClient.pollc                 C   s(   | j j|d |d u r| j  d S d S rN   )r
   closer	   r   r   r   r   rX      s   zKafkaNetClient.closeFc                 C   s.   | j  }|d u r|rt| j j j}|S r   )r
   least_loaded_noderandomchoicer   bootstrap_brokersr   )r   bootstrap_fallbackr   r   r   r   rY      s   
z KafkaNetClient.least_loaded_nodec                    sN    j j }|s|r j j }|s j jd S  fdd|D }t|d S )Nreconnect_backoff_msc                    s   g | ]	} j |jqS r   )r
   connection_delayr   )r5   brokerr   r   r   
<listcomp>   s    z?KafkaNetClient.least_loaded_node_refresh_ms.<locals>.<listcomp>r;   )r
   r   brokersr\   r0   r/   )r   r]   rb   delaysr   r   r   least_loaded_node_refresh_ms   s   z+KafkaNetClient.least_loaded_node_refresh_msc                 C   s   | j |S r   )r
   r_   r   r   r   r   r_      r   zKafkaNetClient.connection_delayc                 C   s   | j   d S r   )r	   wakeupr   r   r   r   re      s   zKafkaNetClient.wakeupc                 C   s(   | j jd u rtd| j jj||dS )Nz?broker_version_data is not available: have we bootstrapped yet?)max_version)r
   broker_version_datar    IllegalStateErrorapi_version)r   	operationrf   r   r   r   ri      s   
zKafkaNetClient.api_version)NN)r&   r   )NrH   )F)__name__
__module____qualname____doc__r   propertyr   r   r   r   r$   r%   r2   r:   rA   rD   rG   rM   rO   rW   r.   rX   rY   rd   r_   re   ri   r   r   r   r   r      s2    

	

	






	r   )loggingrZ   r   r=   kafka.errorserrorsr    kafka.net.managerr   kafka.net.selectorr   	getLoggerrk   logr   r   r   r   r   <module>   s    
