o
    `jf                     @   s   d dl 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 d dlmZ d dlmZ d dlmZmZm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ZdS )    N)Future)KafkaConnectionMetrics)get_sasl_mechanism)ApiVersionsRequest)SaslAuthenticateRequestSaslHandshakeRequestSaslBytesRequest)BrokerVersionData)KafkaProtocol)__version__c                   @   s  e Zd Zi dde dddedddd	d
dddddddddddddddddddddZdTddZedd Zedd Zd d! Z	ed"d# Z
d$d% Zed&d' ZdTd(d)ZdUd*d+ZdTd,d-Zd.d/ Zd0d1 Zd2d3 Zd4d5 Zd6d7 Zd8d9 Zd:d; Zd<d= Zd>d? Zd@dA ZdBdC ZdUdDdEZdFdG ZdHdI ZdUdJdKZdUdLdMZedNdO Z dUdPdQZ!dRdS Z"dS )VKafkaConnection	client_idzkafka-python-client_software_namezkafka-pythonclient_software_version%max_in_flight_requests_per_connection   receive_message_max_bytesi@B request_timeout_msi0u  security_protocol	PLAINTEXTsasl_mechanismNsasl_plain_usernamesasl_plain_passwordsasl_kerberos_namesasl_kerberos_service_namekafkasasl_kerberos_domain_namesasl_oauth_token_providermetricsmetric_group_prefix c                 K   s   t  | j| _| jD ]}||v r|| | j|< q
|| _|| _d | _d | _t | _	t
 | _d| _d| _t | _t | _t | _|| _tj| _d| _t| | _| jd rbt| jd | jd || _nd | _| j| j | j| j d S )NFTr   r   r   )copyDEFAULT_CONFIGconfignode_idnet	transportparsercollectionsdeque_request_buffersetpaused	connectedinitializingr   _init_future_close_futurein_flight_requestsbroker_version_datar   max_version_api_versions_idx_throttle_timeSaslReauthenticator_reauthr   _sensorsadd_errbackfail_in_flight_requestsadd_both)selfr%   r$   r2   configskey r?   R/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/net/connection.py__init__*   s6   




zKafkaConnection.__init__c                 C   s   | j d u rd S | j jS N)r2   broker_versionr<   r?   r?   r@   rC   G   s   
zKafkaConnection.broker_versionc                 C   s   | j  o| j S rB   )r-   r.   rD   r?   r?   r@   closedM   s   zKafkaConnection.closedc                 C   sr   | j rd}n| jsd}n| jrd}nd}| jrd| j  nd}| jd ur(| jnd}d| j | d	| d
| dS )Nr.   disconnectedr,   r-   z
 host=[%s]r    unknownz<KafkaConnection node_id=z broker_version=z (z)>)r.   r-   r,   r&   	host_portrC   r$   )r<   staterH   rC   r?   r?   r@   __str__Q   s   zKafkaConnection.__str__c                 C      | j S rB   )r/   rD   r?   r?   r@   init_future^      zKafkaConnection.init_futurec                 c   s    | j  E d H  | S rB   )rL   	__await__rD   r?   r?   r@   rN   b   s   zKafkaConnection.__await__c                 C   rK   rB   )r0   rD   r?   r?   r@   close_futuref   rM   zKafkaConnection.close_futurec                 C   s^   |d u rt  }|d ur||d  S z|| j W S  ty.   | jd d | _|| j  Y S w )N  r   )time	monotonic_timeout_secsAttributeErrorr#   )r<   now
timeout_msr?   r?   r@   _timeout_atj   s   zKafkaConnection._timeout_atc                 C   s~   t  }| j|d}| js| jjr| j|||f |S | jr*|t	
d| j S | js5|t	dS | j|||d |S )N)rV   zNode paused: zNode not connectedfuture
timeout_at)r   rW   r.   r7   is_reauthenticatingr*   appendr,   failureErrorsNodeNotReadyErrorr-   KafkaConnectionError_send_request)r<   requestr   rY   rZ   r?   r?   r@   send_requestv   s   zKafkaConnection.send_requestc              
      sH   d u rt   jr tdS |jd u r;z	j||_W n tjy: } z |  W  Y d }~S d }~ww t	
 d u rIjdkrV t   S j|}td|| | rj fdd}j| |f n d  jsjj  tjjd krd  S )NrE   )rU   z%s Request %d: %sc                      s     S rB   )_request_timed_outr?   rY   r<   	sent_timerZ   r?   r@   <lambda>   s    z/KafkaConnection._send_request.<locals>.<lambda>r   max_in_flight)r   rE   r]   r^   r`   API_VERSIONr2   api_versionIncompatibleBrokerVersionrQ   rR   rW   KafkaTimeoutErrorr'   rc   logdebugexpect_responser%   call_atr1   r\   successr,   r&   write
send_byteslenr#   pause)r<   rb   rY   rZ   exccorrelation_idtimeout_taskr?   re   r@   ra      sD   



zKafkaConnection._send_requestc                 C   s4   | j r| j  \}}}| j|||d | j sd S d S )NrX   )r*   popleftra   )r<   rb   rY   rZ   r?   r?   r@   send_buffered   s   zKafkaConnection.send_bufferedc                 C   sB   | j s|jrd S || d }td| | | td|  d S )NrP   z6%s: Request timed out after %d ms. Closing connection.zRequest timed out after %d ms)rE   is_donerm   warningcloser^   RequestTimedOutError)r<   rY   sent_atrZ   rV   r?   r?   r@   rd      s
   z"KafkaConnection._request_timed_outc              	   C   s&  | j rtd| t| dS | j|}t|D ]_\}\}}z| j \}}}}	}
W n t	y=   | 
td Y   S w ||krL| 
td  S | j|
 t | d }| jrd| jj| td| ||| | | || qd| jv rt| j| jd k r| d | j  dS )	z# Called when some data is received.z3%s: Ignoring %d bytes received by closed connectionNz-Received response with no in-flight-requests!z$Received unrecognized correlation idrP   z%s: Response %d (%s ms): %srh   r   )rE   rm   rn   rt   r'   receive_bytes	enumerater1   ry   
IndexErrorr}   r^   r`   r%   cancelrQ   rR   r8   request_timerecord_maybe_throttlerq   r,   r#   unpauser7   on_response_processed)r<   data	responsesiresp_correlation_idresponsereq_correlation_idrY   rf   rW   rx   
latency_msr?   r?   r@   data_received   s,   

zKafkaConnection.data_receivedc                 C   s   dS )z Called when the other end calls write_eof() or equivalent.

        If this returns a false value (including None), the transport
        will close itself.  If it returns a true value, closing the
        transport is up to the protocol.
        Fr?   rD   r?   r?   r@   eof_received   s   zKafkaConnection.eof_receivedc                 C   s   d | _ | _d| _| j  |pt }| jjs| j	| | j
jsC|du r4td|  | j
d dS td| | | j
	| dS dS )z Called when the connection is lost or closed.

        The argument is an exception object or None (the latter
        meaning a regular EOF is received or the connection was
        aborted or closed).
        FNz%s: Connection closedz%s: Connection lost: %s)r-   r.   r&   r7   r   r^   r`   r/   r{   r]   r0   rm   inforq   error)r<   rv   r   r?   r?   r@   connection_lost   s   
zKafkaConnection.connection_lostc                 C   s~   | j std|pt }| jr | j \}}}|| | js| jr=| j \}}}}}| j	| || | js#d S d S )Nz4Connection must be closed to fail in flight requests)
rE   RuntimeErrorr^   	Cancelledr*   ry   r]   r1   r%   r   )r<   r   _rY   rx   r?   r?   r@   r:      s   

z'KafkaConnection.fail_in_flight_requestsc                 C   s   | j rtd|| _| j | kr| j|  d| _| j  zd| jg| j	 dd R  }W n t
yB   td d}Y nw t| jd | jd	 |d
| _td|  dS )z Called when a connection is made.

        The argument is the transport representing the pipe connection.
        To receive data, wait for data_received() calls.
        When the connection is closed, connection_lost() is called.
        z Connection closed during connectTznode=%s[%s:%s]r      z%Failed to build connection log_prefixr    r   r   )r   r   identz%s: Connection madeN)rE   r^   r`   r&   get_protocolset_protocolr.   resume_readingr$   getPeer	Exceptionrm   	exceptionr
   r#   r'   rn   )r<   r&   
log_prefixr?   r?   r@   connection_made  s&   

$
zKafkaConnection.connection_madec                 C   s   | j | d S rB   )r,   add)r<   vr?   r?   r@   ru   "  s   zKafkaConnection.pausec                 C   sf   z| j | W n
 ty   Y d S w | j s+| jr-| jr/| j }|r1| j| d S d S d S d S d S rB   )r,   removeKeyErrorr'   r&   rs   rr   )r<   r   to_sendr?   r?   r@   r   %  s   
zKafkaConnection.unpausec                 C      |  d dS )a   Called when the transport's buffer goes over the high-water mark.

        Pause and resume calls are paired -- pause_writing() is called
        once when the buffer goes strictly over the high-water mark
        (even if subsequent writes increases the buffer size even
        more), and eventually resume_writing() is called once when the
        buffer size reaches the low-water mark.

        Note that if the buffer size equals the high-water mark,
        pause_writing() is not called -- it must go strictly over.
        Conversely, resume_writing() is called when the buffer size is
        equal or lower than the low-water mark.  These end conditions
        are important to ensure that things go as expected when either
        mark is zero.

        NOTE: This is the only Protocol callback that is not called
        through EventLoop.call_soon() -- if it were, it would have no
        effect when it's most needed (when the app keeps writing
        without yielding until pause_writing() is called).
        bufferN)ru   rD   r?   r?   r@   pause_writing0  s   zKafkaConnection.pause_writingc                 C   r   )zD Called when the transport's buffer drains below the low-water mark.r   N)r   rD   r?   r?   r@   resume_writingG     zKafkaConnection.resume_writingc                 C   sN   |d u r| j jst }| js| | d S |r | j| d S | j  d S rB   )r/   r{   r^   r`   r&   r   abortr}   )r<   r   r?   r?   r@   r}   K  s   
zKafkaConnection.closec                 C   s   t |dd}| jr| jj| |sd S | jd ur;| jdkr;t |d  }|| jkr;|| _| j	|| j
 | d td| |jj| d S )Nthrottle_time_msr   )r   r   rP   throttlez"%s: %s throttled by broker (%d ms))getattrr8   throttle_timer   rC   rQ   rR   r5   r%   rp   _maybe_unthrottleru   rm   r|   	__class____name__)r<   r   r   r   r?   r?   r@   r   V  s   

zKafkaConnection._maybe_throttlec                 C   s&   t  | jkrd| _| d d S d S )Nr   r   )rQ   rR   r5   r   rD   r?   r?   r@   r   g  s   z!KafkaConnection._maybe_unthrottlec              
      s   |d u r	|   }z| |I d H  | jr| |I d H  W n ty6 } z| | W Y d }~d S d }~ww |   td|  d S )Nz%s: Connected)	rW   _get_api_versionssasl_enabled_sasl_authenticater   r}   _init_completerm   r   )r<   rZ   r   r?   r?   r@   
initializel  s   zKafkaConnection.initializec           	         s  |d u r	|   }| jd ur+z	| jt| _W n tjy*   td| | j	 Y d S w |t
 kry| j}t|| jd | jd d}| j||dI d H }t|j}|tju rWn'|tju rv|jD ]}|j|jkrqt| j|j| _ nq_d| _q+| tddd	 |jD }t|d
}| jd u rtd| |j || _d S | j|krtd| |j| jj || _d S | jd ur| j|k rtd| |j| jj d S d S d S )Nz7%s: Using pre-configured api_version %s for ApiVersionsr   r   )versionr   r   rZ   r   z Timeout during ApiVersions checkc                 S   s   i | ]
}|j |j|jfqS r?   )api_keymin_versionr3   ).0rj   r?   r?   r@   
<dictcomp>  s    z5KafkaConnection._get_api_versions.<locals>.<dictcomp>)api_versionsz#%s: Broker version identified as %szA%s: Broker version identified as %s (lower than user-supplied %s)zA%s: Broker version identified as %s; clamping to user-supplied %s)rW   r2   rj   r   r4   r^   rk   rm   rn   rC   rQ   rR   r#   ra   for_code
error_codeNoErrorUnsupportedVersionErrorapi_keysr   API_KEYminr3   rl   r	   r   broker_version_str)	r<   rZ   r   rb   r   
error_typerj   r   bvdr?   r?   r@   r   y  sX   









z!KafkaConnection._get_api_versionsc                 C   s   | j d dv S )Nr   )SASL_PLAINTEXTSASL_SSL)r#   rD   r?   r?   r@   r     r   zKafkaConnection.sasl_enabledc                    s  |d u r	|   }t| jd dd}| j||dI d H }t|j}|tjur2t	d| |j
 | | jd |jvrGtd| jd |jf |j}| jjrR| jjn| j d }t| jd dd|i| j}d }| s|t kr| }	|dkrt|	}
nt|	}
| j|
|dI d H }t|j}|tjurtd	|j
|jf |dkr| rn||j | s|t kstt |krtd
| std| jd  |dkr| j|j t d	| |!  d S )Nr      )	mechanismr3   r   z%s: SaslHandshake failed: %szGKafka broker does not support %s sasl mechanism. Enabled mechanisms: %sr   hostz%s: %szSASL Authentication timed outz"Failed to authenticate via SASL %sr?   )"rW   r   r#   ra   r^   r   r   r   rm   r   r   
mechanismsUnsupportedSaslMechanismErrorri   r&   r   r   r   r{   rQ   rR   
auth_bytesr   r   SaslAuthenticationFailedErrorerror_messagereceiverl   is_authenticatedr7   session_updatedsession_lifetime_msr   auth_details)r<   rZ   rb   r   r   r   	sasl_hostr   auth_responsetokenauth_requestr?   r?   r@   r     sd   



z"KafkaConnection._sasl_authenticatec                 C   s8   | j rd| _ d| _|   | jd | j  d S d S )NFT)r.   r-   rz   r/   rq   r7   schedulerD   r?   r?   r@   r     s   zKafkaConnection._init_complete)NNrB   )#r   
__module____qualname__r   r"   rA   propertyrC   rE   rJ   rL   rN   rO   rW   rc   ra   rz   rd   r   r   r   r:   r   ru   r   r   r   r}   r   r   r   r   r   r   r   r?   r?   r?   r@   r      s    	








*
	


/

8r   c                   @   s`   e Zd ZdZdd Zedd Zedd Zdd	 Zd
d Z	dd Z
dd Zdd Zdd ZdS )r6   a  KIP-368 SASL re-authentication state and scheduling for a single
    KafkaConnection. Owns the per-connection re-auth lifecycle so the
    connection doesn't have to carry the related attributes and coroutines
    inline. The connection plugs this in at five points:

      - after each successful SASL auth                  -> session_updated()
      - after init completes                             -> schedule()
      - when send_request needs to gate the public API   -> is_reauthenticating
      - on every response popped from in_flight_requests -> on_response_processed()
      - on connection_lost                               -> cancel()
    c                 C   s(   || _ d| _d | _d | _d| _d | _d S )Nr   F)_connr   authenticated_at_task_reauthenticating_drain_future)r<   connr?   r?   r@   rA     s   
zSaslReauthenticator.__init__c                 C   rK   rB   )r   rD   r?   r?   r@   r[     rM   z'SaslReauthenticator.is_reauthenticatingc                 C   rK   )zEThe scheduled re-auth task, or None. Exposed for tests/observability.)r   rD   r?   r?   r@   task  s   zSaslReauthenticator.taskc                 C   sJ   |pd| _ | j dk rd| _ nd| j   k rdkrn nd| _ t | _dS )zCapture broker-advertised session lifetime after each successful
        auth round (initial and subsequent re-auths). Clamp negative values to 0,
        and require minimum non-zero lifetime of 1sec (1000).r   rP   N)r   rQ   rR   r   )r<   r   r?   r?   r@   r     s   

z#SaslReauthenticator.session_updatedc                 C   sX   | j jr| js	dS tdd}| j| d }td| j || j | j j|| j	| _
dS )a  Schedule the next re-auth before the lifetime elapses. Jittered to
        85-95% of the lifetime to avoid synchronised re-auth storms across
        many connections (Apache Java semantics). No-op when SASL is disabled
        or the broker advertised lifetime=0.
        Ng333333?gffffff?rP   zG%s: Scheduling SASL re-authentication in %.3fs (session_lifetime_ms=%d))r   r   r   randomuniformrm   rn   r%   
call_later_runr   )r<   pctdelayr?   r?   r@   r     s   
zSaslReauthenticator.schedulec                 C   sR   | j dur| jj| j  d| _ | jdur!| jjs!| jt  d| _d| _	dS )zvCancel any pending re-auth and fail the drain awaiter if present.
        Called from KafkaConnection.connection_lost.NF)
r   r   r%   r   r   r{   r]   r^   r`   r   rD   r?   r?   r@   r   $  s   

zSaslReauthenticator.cancelc                 C   s@   | j r| jdur| jjs| jjs| jd dS dS dS dS dS )zWake the drain awaiter once in_flight_requests clears during reauth.
        Called from KafkaConnection.data_received after each pop.N)r   r   r   r1   r{   rq   rD   r?   r?   r@   r   /  s   
z)SaslReauthenticator.on_response_processedc              
      s   d | _ | jjr
d S z
|  I d H  W d S  tyC } z#td| j| t|tr+|nt	
t|}| j| W Y d }~d S d }~ww )Nz%%s: SASL re-authentication failed: %s)r   r   rE   
_do_reauthBaseExceptionrm   r|   
isinstancer   r^   r   strr}   )r<   rv   errr?   r?   r@   r   8  s   zSaslReauthenticator._runc                    s   d| _ zF| jjr$| jjs$t | _| jjsn| jI d H  | jjr$| jjrd | _| jjr4W d| _ d | _d S td| j | j I d H  W d| _ d | _nd| _ d | _w | jjrXd S | j	  | 
  d S )NTFz$%s: Beginning SASL re-authentication)r   r   r1   rE   r   r   rm   rn   r   rz   r   rD   r?   r?   r@   r   F  s0   

zSaslReauthenticator._do_reauthN)r   r   r   __doc__rA   r   r[   r   r   r   r   r   r   r   r?   r?   r?   r@   r6     s    

	r6   ) r(   r!   loggingr   structrQ   kafka.errorserrorsr^   kafka.futurer   kafka.net.metricsr   kafka.net.saslr   kafka.protocol.metadatar   kafka.protocol.saslr   r   r   "kafka.protocol.broker_version_datar	   kafka.protocol.parserr
   kafka.versionr   	getLoggerr   rm   r   r6   r?   r?   r?   r@   <module>   s*    
   Z