o
    `jr                     @   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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dd ZdZdd Zd	d
 ZG dd dZG dd dejZG dd dZG dd dZdS )    N)Future)__version__c                 O   s   t j d| g|R i | d S )N   )log)msgargskwargs r	   P/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/net/selector.py	log_trace   s   r   i  c                 c   s    |  V  d S Nr	   )callbackr	   r	   r
   yield_callback   s   r   c                 C   s\   t | s
t | r| S t | st | r| S t | s"t | r&t| S tdt	|  )Nz$Generator or coroutine not found: %s)
inspectisgeneratoriscoroutineisgeneratorfunctioniscoroutinefunction
isfunctionismethodr   	TypeErrortype)
maybe_coror	   r	   r
   _initialize_coro   s   r   c                   @   s   e Zd Zdd Zdd ZdS )KernelEventc                 G   s   || _ || _d S r   )methodr   )selfr   r   r	   r	   r
   __init__,   s   
zKernelEvent.__init__c                 c   s    | V S r   r	   r   r	   r	   r
   	__await__0   s   zKernelEvent.__await__N)__name__
__module____qualname__r   r   r	   r	   r	   r
   r   +   s    r   c                   @   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 )	TaskStatecreated	scheduledunscheduledreadyrunningwait_iowait_futuredone	cancelledN)r    r!   r"   CREATED	SCHEDULEDUNSCHEDULEDREADYRUNNINGWAIT_IOWAIT_FUTUREDONE	CANCELLEDr	   r	   r	   r
   r#   4   s    r#   c                   @   sb   e Zd Zdd Zdd ZdddZdd	 Zd
d Zdd Ze	dd Z
e	dd Ze	dd ZdS )Taskc                 C   s,   t |d f| _d | _d | _d | _tj| _d S r   )r   _stack_res_excscheduled_atr#   r-   stater   coror	   r	   r
   r   A   s
   zTask.__init__c                 C   s   t | t |k S r   )id)r   otherr	   r	   r
   __lt__H   s   zTask.__lt__Nc              
   C   sv  | j rtd| jd ur| jd }| _d }nd }d }	 | jd }t|r9t|s9t|s9| }|| jd f| _z/|rB||}n|	|}t
|ttfrQ|W S t|s`t|s`t|rg| | d }W nO ty } z| jd | _| jstj| _|j| _ |j}d }W Y d }~n-d }~w ty } z| jd | _| jstj| _|| _ d }|}W Y d }~nd }~ww d }q)NTask is already done!Tr      )is_doneRuntimeErrorr9   r7   callabler   r   r   throwsend
isinstancer   r   r   
push_stackStopIterationr#   r4   r;   valuer8   BaseException)r   argexcretr=   finaler	   r	   r
   __call__N   sV   




zTask.__call__c                 C   s   t || jf| _d S r   )r   r7   r<   r	   r	   r
   rI         zTask.push_stackc                 C   s<   | j rtdt|tstd| jd urtd|| _d S )NrA   zexc is not a BaseExceptionzTask exception is already set!)rC   rD   rH   rL   r   r9   r   rN   r	   r	   r
   
inject_exc   s   


zTask.inject_excc                 C   s   | j rd S | jtjusJ | j}|r7|\}}t|s t|r5z|  W n t	y4   t
d Y nw |sd | _tj| _t | _d S )Nz*Error closing coroutine for cancelled task)rC   r;   r#   r1   r7   r   r   r   close	Exceptionr   	exceptionr5   Errors	Cancelledr9   )r   stackr=   r	   r	   r
   rV      s    z
Task.closec                 C   s
   | j d u S r   )r7   r   r	   r	   r
   rC      s   
zTask.is_donec                 C      | j std| jS NzTask not complete!)rC   rD   r8   r   r	   r	   r
   result      zTask.resultc                 C   r\   r]   )rC   rD   r9   r   r	   r	   r
   rX      r_   zTask.exceptionr   )r    r!   r"   r   r@   rR   rI   rU   rV   propertyrC   r^   rX   r	   r	   r	   r
   r6   @   s    
3	

r6   c                   @   sp  e Zd Zde ejdddZdd Zdd Zd	d
 Z	dd Z
dVddZdd Zdd ZdWd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.d/ ZdVd0d1ZdVd2d3ZdVd4d5ZdVd6d7Zd8d9 Z d:d; Z!d<d= Z"d>d? Z#d@dA Z$dBdC Z%dDdE Z&dFdG Z'dHdI Z(dXdJdKZ)dVdLdMZ*dNdO Z+dPdQ Z,dRdS Z-dTdU Z.dS )YNetworkSelectorzkafka-python-g?F)	client_idselectorslow_task_threshold_secsraise_on_slow_taskc                 K   s   t  | j| _| jD ]}||v r|| | j|< q
t | _d | _d| _d | _d| _	| jd  | _
g | _t | _t | _d | _t \| _| _| jd | jd | j
| jtjd d | _i | _t | _d S )NFrc   NN)copyDEFAULT_CONFIGconfig	threadingLock
_poll_lock_poll_owner_closed
_exception_stop	_selector
_scheduledcollectionsdeque_readyset_pending_tasks_currentsocket
socketpair	_wakeup_r	_wakeup_wsetblockingregister	selectors
EVENT_READ
_io_thread_pending_waiters_pending_waiters_lock)r   configskeyr	   r	   r
   r      s,   


	zNetworkSelector.__init__c                 C   s$   dt | jt | jt | j f S )Nz2<NetworkSelector ready=%d scheduled=%d waiting=%d>)lenru   rr   rq   get_mapr   r	   r	   r
   __str__   s   $zNetworkSelector.__str__c              
   C   s   d| _ td| jd  z| j s|   | j r|   W n ty: } ztd| jd  || _| 	|  d}~ww td| jd | j  dS )zRun the event loop until stop() is called. Intended to be driven by
        a dedicated IO thread. Wake-ups from other threads must go through
        call_soon_threadsafe() so the select() loop returns promptly.FzIO loop starting (client_id=%s)rb   zIO loop crashed (client_id=%s)Nz.IO loop exited cleanly (client_id=%s, stop=%s))
rp   r   infori   
_poll_oncedrainrL   rX   ro   _fail_pending_waitersrT   r	   r	   r
   run_forever   s"   
zNetworkSelector.run_foreverc                 C   s<   | j durdS tj| jd| jd  dd}|| _ |  dS )z>Spawn a daemon IO thread that owns the event loop. Idempotent.Nzkafka-io-%srb   T)targetnamedaemon)r   rj   Threadr   ri   start)r   tr	   r	   r
   r      s   
zNetworkSelector.startNc                 C   sX   | j s| jdu r
dS d| _ |   | j|dur|d nd d| _| td dS )aC  Signal run_forever() to exit and join the IO thread.

        Blocks the caller until the IO thread terminates (or ``timeout_ms``
        elapses). Pending cross-thread ``run()`` waiters are failed with
        KafkaConnectionError. Idempotent; safe to call from any thread
        other than the IO thread itself.
        NT  zEvent loop stopped)rp   r   wakeupjoinr   rY   KafkaConnectionError)r   
timeout_msr	   r	   r
   stop  s   zNetworkSelector.stopc                 C   s`   | j  t| j }| j  W d    n1 sw   Y  |D ]\}}||d< |  q!d S )NrX   )r   listr   itemsclearrv   )r   rN   waiterseventr;   r	   r	   r
   r     s   
z%NetworkSelector._fail_pending_waitersc                    s   j rtdjdu r&jg R  }j|d |jdur#|j|jS t ju r1tdj	r8j	dt
 ddd fdd}j j< W d   n1 s^w   Y  |   d durvd d	 S )
a  Schedules coro on the event loop, blocks until complete, returns value or raises.

        If an IO thread is running (via start()), the caller thread blocks on
        a cross-thread Event while the coroutine runs on the IO thread. Safe
        to call concurrently from multiple caller threads.

        If no IO thread is running, falls back to driving the loop on the
        caller thread (legacy behavior).
        NetworkSelector closed!N)futurea$  Cannot block on net.run() from the IO thread itself. This typically happens when a synchronous rebalance listener (or another IO-thread callback) calls a blocking consumer/admin API. Use AsyncConsumerRebalanceListener and await the async variant, or move the blocking work to a worker thread.)rK   rX   c                     s   z]zj g R  I d H d< W n+ ty= }  zd d u r%| d< nt| ts3tdd |  W Y d } ~ nd } ~ ww W j jd  W d    n1 sTw   Y  	  d S j jd  W d    n1 stw   Y  	  w )NrK   rX   z:During exception %s, caught additional error %s (ignoring))
_invokerL   rH   GeneratorExitr   warningr   r   poprv   rN   r   r=   r   r   r;   r	   r
   waiter:  s&    


z#NetworkSelector.run.<locals>.waiterrX   rK   )rn   rD   r   call_soon_with_futurepollrX   rK   rj   current_threadro   Eventr   r   call_soon_threadsafewait)r   r=   r   r   r   r	   r   r
   run  s2   




zNetworkSelector.runc                 C   s8   | j s|r| jr|   | j s|r| jsd S d S d S d S r   )ru   rr   r   )r   r%   r	   r	   r
   r   O  s    zNetworkSelector.drainc                 C   sP   | j rtdt|tst|}||_tj|_t	| j
||f | j| |S Nr   )rn   rD   rH   r6   r:   r#   r.   r;   heapqheappushrr   rw   addr   whentaskr	   r	   r
   call_atS  s   
zNetworkSelector.call_atc                 C   s*   t |ts	t|}| t | | |S r   )rH   r6   r   time	monotonic)r   delayr   r	   r	   r
   
call_later^  s   
zNetworkSelector.call_laterc                 C   s   | j | tj|_d S r   )ru   appendr#   r0   r;   r   r   r	   r	   r
   _add_ready_taskd  s   zNetworkSelector._add_ready_taskc                 C   s&   |j std| j| tj|_d S )NzTask is not done yet!)rC   rD   rw   discardr#   r4   r;   r   r	   r	   r
   
_task_doneh  s   zNetworkSelector._task_donec                 C   s,   t |ts	t|}| | | j| |S r   )rH   r6   r   rw   r   r   r	   r	   r
   	call_soonn  s
   

zNetworkSelector.call_soonc                 C   s2   | j r| j d | jrtd| |}|   |S r   )ro   rn   rD   r   r   )r   r   r   r	   r	   r
   r   u  s   
z$NetworkSelector.call_soon_threadsafec                    s<   t dr rtdt  fdd}| S )Nr   z(initiated coroutine does not accept argsc               
      sX   z jg R  I d H  W d S  ty+ }  z|  W Y d } ~ d S d } ~ ww r   )successr   rL   failurer   r   r=   r   r   r	   r
   wrapper  s   $z6NetworkSelector.call_soon_with_future.<locals>.wrapper)hasattr
ValueErrorr   r   )r   r=   r   r   r	   r   r
   r   ~  s   

z%NetworkSelector.call_soon_with_futurec                    sz   t |r|| I dH }nt|dr|I dH }n|| }t |s't|dr,|I dH }t|tr;|I dH }t|ts1|S )zInvoke coro/awaitable/function and fully resolve the result.

        If the result is itself a Future (e.g. send() returning an unresolved
        Future), it is awaited so callers receive the resolved value.
        Nr   )r   r   r   r   rH   r   )r   r=   r   r^   r	   r	   r
   r     s   





zNetworkSelector._invokec                 C   sf   |j tju sJ |jd usJ z| j|j|f W n	 ty#   Y nw t| j d |_tj	|_ d S r   )
r;   r#   r.   r:   rr   remover   r   heapifyr/   r   r	   r	   r
   _unschedule  s   zNetworkSelector._unschedulec                 C   s|   |j tjtjfv rd S |j tju r|| ju sJ tj| j_ d S |j tju r+| | n|j tju r2	 | j	
| |  d S r   )r;   r#   r4   r5   r1   rx   r.   r   r2   rw   r   rV   r   r	   r	   r
   cancel  s   
zNetworkSelector.cancelc                 C   s&   |j tju r| | | || |S r   )r;   r#   r.   r   r   r   r	   r	   r
   
reschedule  s   
zNetworkSelector.reschedulec                 C   s
   t d|S )N_sleepr   r   r   r	   r	   r
   sleep  s   
zNetworkSelector.sleepc                 C   s   |  || j d S r   )r   rx   r   r	   r	   r
   r        zNetworkSelector._sleepc                 C      t d||S )N_wait_writer   r   fileobj
timeout_atr	   r	   r
   
wait_write     zNetworkSelector.wait_writec                 C      |  |tj| d S r   )_wait_ior   EVENT_WRITEr   r	   r	   r
   r     rS   zNetworkSelector._wait_writec                 C   r   )N
_wait_readr   r   r	   r	   r
   	wait_read  r   zNetworkSelector.wait_readc                 C   r   r   )r   r   r   r   r	   r	   r
   r     rS   zNetworkSelector._wait_readc                    s|   j   d  fdd}| }t| | tj_|d u s,jr.d S fdd}||d S )Nc                
   3   sZ    zd V  W d urj s   d S d ur&j s&   w r   )rC   r   unregister_eventr	   )r   r   r   timerr	   r
   io_guard  s   

z*NetworkSelector._wait_io.<locals>.io_guardc                      s,   d j rd S td   d S )NzI/O wait timed out)rC   rU   rY   KafkaTimeoutErrorr   r	   )r   	suspendedr   r	   r
   
on_timeout  s
   z,NetworkSelector._wait_io.<locals>.on_timeout)	rx   register_eventnextrI   r#   r2   r;   rn   r   )r   r   r   r   r   guardr   r	   )r   r   r   r   r   r
   r     s   
zNetworkSelector._wait_ioc                 C   sh   | j r.| j d d t kr2t| j \}}d |_| | | j r0| j d d t ksd S d S d S d S Nr   )rr   r   r   r   heappopr:   r   )r   _r   r	   r	   r
   _schedule_tasks  s
   
,zNetworkSelector._schedule_tasksc                 C   s*   z
| j d d | W S  ty   Y d S w r   )rr   
IndexError)r   nowr	   r	   r
   _next_scheduled_timeout  s
   z'NetworkSelector._next_scheduled_timeoutc              	   C   s   t d||| t|tst|}zD| j|}|j\}}|tjkr,|r,td|||f |tj	kr<|r<td|||f | j
||j|B |tjkrM||fn||f W d S  tyq   | j|||tjkri|d fnd |f Y d S w )Nznet.register_event: %s, %s, %sz<EVENT_READ already registered for fileobj %s by %s (new: %s)z=EVENT_WRITE already registered for fileobj %s by %s (new: %s))r   rH   r6   rq   get_keydatar   r   rD   r   modifyeventsKeyErrorr~   )r   r   r   r   r   readerwriterr	   r	   r
   r     s   

2,zNetworkSelector.register_eventc              	   C   s   t d|| z2| j|}|j\}}|j| @ }|s#| j| W d S | j|||tjkr1d |fn|d f W d S  t	t
fyD   Y d S w )Nznet.unregister_event: %s, %s)r   rq   r   r   r   
unregisterr   r   r   r   r   )r   r   r   r   r   r   r   r	   r	   r
   r     s   
,z NetworkSelector.unregister_eventc                 C   r   r   )r   r   r   r   r   r   r	   r	   r
   
add_reader$  rS   zNetworkSelector.add_readerc                 C      |  |tj d S r   )r   r   r   r   r   r	   r	   r
   remove_reader'  r   zNetworkSelector.remove_readerc                 C   r   r   )r   r   r   r   r	   r	   r
   
add_writer*  rS   zNetworkSelector.add_writerc                 C   r   r   )r   r   r   r   r	   r	   r
   remove_writer-  r   zNetworkSelector.remove_writerc                 C   s   t d t }|d ur|d nd }|d ur|jrd}	 | | |d u s(|jr)n|d ur<|d t |  }|dkr<nqt d d S )Nzpoll: enterr   r   Tz
poll: exit)r   r   r   rC   r   )r   r   r   start_atinner_timeoutr	   r	   r
   r   0  s   
zNetworkSelector.pollc                    s|  t d  jjdds jt u rtdtdt  _z jr'd}n t	
 }|d ur=|d ur;t||n|}|d urO|tkrHt}n|dk rNd}n j sVd} j|}t dt|  |     jd }t j}t|D ]#} j  _ jjtjtjfv rqztj j_|rt	
 nd }zzt d	 j   }W n$ ty     j Y n ty   t d
 j   j Y nsw  jjtju r j!" j  j#  n^t$|t%r!t d|j& zt' |j&|j(  W nF ty  }	 zt d|j&|	 j  j)|	  * j W Y d }	~	n#d }	~	ww t$|t+r9|, jf fdd	 tj- j_ntd| W  jd urZ jjtju rZt.d j tj/ j_n jd uru jjtju rut.d j tj/ j_w |rt	
 | }
|
|krd j|
|f } jd rd  _t|t.| qzd  _W d  _ j0  t d d S d  _ j0  t d w )Nz_poll_once: enterF)blockingzRecursive access to net.poll!zConcurrent access to net.poll!r   z_poll_once: %d ready_eventsrd   zCalling task %szUnhandled exception in task %s:zkernel event %sz,kernel event %s raised %r; injecting into %sc                    s
     |S r   )r   )r   r   r   r	   r
   <lambda>  s   
 z,NetworkSelector._poll_once.<locals>.<lambda>zUnhandled event type: %sz8Task %s left RUNNING after step; demoting to UNSCHEDULEDzTask %r ran for %.3fs (>%.3fs threshold). It is blocking the event loop -- likely a tight sync loop inside a coroutine. Other pollers will time out.re   z_poll_once: exit)1r   rl   acquirerm   rj   r   rD   ru   r   r   r   minMAX_TIMEOUTrq   r   selectr   _process_eventsr   ri   rangepopleftrx   r;   r#   r4   r5   r1   rJ   r   rL   r   rX   rw   r   rV   rH   r   r   getattrr   rU   r   r   add_bothr3   r   r/   release)r   timeoutscheduled_timeoutready_events	thresholdni
step_startr   rQ   elapsedr   r	   r   r
   r   @  s   













zNetworkSelector._poll_oncec              	   C   s,   z	| j d W d S  ttfy   Y d S w )N    )r|   rG   BlockingIOErrorOSErrorr   r	   r	   r
   r     s
   zNetworkSelector.wakeupc              	   C   s   | j | jfD ]#}z| j| W n	 ty   Y nw z|  W q ty)   Y qw t \| _ | _| j d | jd | j	| j t
jd d S )NFrf   )r{   r|   rq   r   rW   rV   ry   rz   r}   r~   r   r   )r   sr	   r	   r
   _rebuild_wakeup_socketpair  s   z*NetworkSelector._rebuild_wakeup_socketpairc              	   C   s   | j rd S d| _ | jd ur|   |   t| jD ]}| | q| j| jfD ]#}z| j	
| W n	 ty;   Y nw z|  W q( tyK   Y q(w | j	  d S )NT)rn   r   r   r   r   rw   r   r{   r|   rq   r   rW   rV   )r   r   r  r	   r	   r
   rV     s(   
zNetworkSelector.closec           	      C   s   |D ]r\}}|j \}}|j}|| ju rLz| jd}|s%td |   W n$ ty/   Y n tyJ } ztd| |   W Y d }~nd }~ww q|t	j
@ r`|d ur[| | ntd |t	j@ rt|d uro| | qtd qd S )Ni   z)Wakeup socket returned empty. Rebuilding.z,Error reading wakeup socket: %s. Rebuilding.z*Selector got WRITE event without writer...z)Selector got READ event without reader...)r   r   r{   recvr   r   r  r  rW   r   r   r   r   )	r   
event_listr   r   r   r   r   r   rQ   r	   r	   r
   r    s8   






zNetworkSelector._process_eventsr   )Frf   )/r    r!   r"   r   r   DefaultSelectorrh   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r  rV   r  r	   r	   r	   r
   ra      s\    &


3	



%

era   )rs   rg   enumr   loggingr   r   ry   rj   r   kafka.errorserrorsrY   kafka.futurer   kafka.versionr   	getLoggerr    r   r   r  r   r   r   Enumr#   r6   ra   r	   r	   r	   r
   <module>   s,    
	n