o
    `j@                     @   s@   d dl Z d dlZd dlZd dlmZ eeZG dd dZdS )    N)RetriableErrorc                   @   st   e Zd ZdZdZdd Zdd Zdd Zd	d
 Zdd Z	dd Z
dd Zdd Zdd Zdd Zdd Zdd ZdS )Future)is_donevalue	exception
_callbacks	_errbacks_lockFc                 C   s,   d| _ d | _d | _g | _g | _t | _d S )NF)r   r   r   r   r   	threadingLockr	   self r   J/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/future.py__init__   s   zFuture.__init__c                 C   s   | j o| jd u S Nr   r   r   r   r   r   	succeeded      zFuture.succeededc                 C   s   | j o| jd uS r   r   r   r   r   r   failed   r   zFuture.failedc                 C   s   | j ot| jtS r   )r   
isinstancer   r   r   r   r   r   	retriable   s   zFuture.retriablec                 C   s   | j }|  || _d| _| j}d | _d | _|  |rE| j}|D ]#}z|| W q! tyD } zt	
d |r:|W Y d }~q!d }~ww | S )NTError processing callback)r	   acquirer   r   r   r   releaseerror_on_callbacks	Exceptionlogr   )r   r   lock	callbacksr   fer   r   r   success   s*   
zFuture.successc                 C   s   t |t ur|n| }t|tstd| j}|  || _d| _| j}d | _	d | _|
  |rY| j}|D ]#}z|| W q5 tyX } ztd |rN|W Y d }~q5d }~ww | S )Nz"future failed without an exceptionTError processing errback)typer   BaseException	TypeErrorr	   r   r   r   r   r   r   r   r   r   )r   r!   r   r   errbacksr   r    errr   r   r   failure:   s0   

zFuture.failurec              
   O   s   |s|rt j|g|R i |}| j}|  | js&| j| |  | S |  | jd u rUz|| j	 W | S  t
yT } ztd | jrI|W Y d }~| S d }~ww | S )Nr   )	functoolspartialr	   r   r   r   appendr   r   r   r   r   r   r   r    argskwargsr   r!   r   r   r   add_callbackQ   ,   


zFuture.add_callbackc              
   O   s   |s|rt j|g|R i |}| j}|  | js&| j| |  | S |  | jd urUz|| j W | S  t	yT } zt
d | jrI|W Y d }~| S d }~ww | S )Nr#   )r*   r+   r	   r   r   r   r,   r   r   r   r   r   r-   r   r   r   add_errbackd   r1   zFuture.add_errbackc              
   C   s   | j }|  | js| j| | j| |  | S |  | jdu rKz|| j W | S  t	yJ } zt
d | jr?|W Y d}~| S d}~ww z|| j W | S  t	yp } zt
d | jre|W Y d}~| S d}~ww )a  Register a (callback, errback) pair under a single lock acquire.

        Fast path for call sites that always register both a plain callback
        and errback with no ``*args``/``**kwargs``. Used on the producer hot
        path (``FutureRecordMetadata`` -> ``FutureProduceResult``) to halve
        the per-record lock-acquire count vs. calling ``add_callback()`` +
        ``add_errback()`` separately.
        Nr   r#   )r	   r   r   r   r,   r   r   r   r   r   r   r   )r   cbebr   r!   r   r   r   
_add_cb_ebw   s>   	


	

zFuture._add_cb_ebc                 O   s4   | j |g|R i | | j|g|R i | | S r   )r0   r2   )r   r    r.   r/   r   r   r   add_both   s   zFuture.add_bothc                 C   s   |  |j|j | S r   )r5   r"   r)   )r   futurer   r   r   chain   s   zFuture.chainc                 c   s     | j s| V  | jr| j| jS r   )r   r   r   r   r   r   r   	__await__   s   zFuture.__await__N)__name__
__module____qualname__	__slots__r   r   r   r   r   r"   r)   r0   r2   r5   r6   r8   r9   r   r   r   r   r   
   s    !r   )	r*   loggingr
   kafka.errorsr   	getLoggerr:   r   r   r   r   r   r   <module>   s    
