
Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­
<!DOCTYPE html>
<html>
U
    ¡ê,aª-  ã                   @   sÐ   d ddg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	m
Z
 ddlZddlmZ ddlmZ ejjZdd	lmZmZmZmZmZ G d
d „ d eƒZeƒ ZG dd„ deƒZG dd„ deƒZdS )ÚQueueÚSimpleQueueÚJoinableQueueé    N)ÚEmptyÚFullé   )Ú
connection)Úcontext)ÚdebugÚinfoÚFinalizeÚregister_after_forkÚ
is_exitingc                   @   sº   e Zd Zd*dd„Zdd„ Zdd„ Zdd	„ Zd+dd„Zd,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ed"d#„ ƒZed$d%„ ƒZed&d'„ ƒZed(d)„ ƒZdS )-r   r   c                C   s’   |dkrddl m} || _tjdd�\| _| _| ¡ | _t	 
¡ | _tjdkrTd | _n
| ¡ | _| |¡| _d| _|  ¡  tjdkrŽt| tjƒ d S )Nr   r   )ÚSEM_VALUE_MAXF©ZduplexÚwin32)Zsynchronizer   Ú_maxsizer   ÚPipeÚ_readerÚ_writerÚLockÚ_rlockÚosÚgetpidÚ_opidÚsysÚplatformÚ_wlockZBoundedSemaphoreÚ_semÚ_ignore_epipeÚ_after_forkr   r   ©ÚselfÚmaxsizeÚctx© r%   ú;/opt/alt/python38/lib64/python3.8/multiprocessing/queues.pyÚ__init__$   s    




zQueue.__init__c                 C   s.   t  | ¡ | j| j| j| j| j| j| j| j	fS ©N)
r	   Úassert_spawningr   r   r   r   r   r   r   r   ©r"   r%   r%   r&   Ú__getstate__9   s    
   ÿzQueue.__getstate__c              	   C   s0   |\| _ | _| _| _| _| _| _| _|  ¡  d S r(   )	r   r   r   r   r   r   r   r   r    ©r"   Ústater%   r%   r&   Ú__setstate__>   s    ÿ   zQueue.__setstate__c                 C   sb   t dƒ t t ¡ ¡| _t ¡ | _d | _d | _	d| _
d| _d | _| jj| _| jj| _| jj| _d S )NzQueue._after_fork()F)r
   Ú	threadingÚ	Conditionr   Ú	_notemptyÚcollectionsÚdequeÚ_bufferÚ_threadÚ_jointhreadÚ_joincancelledÚ_closedÚ_closer   Ú
send_bytesÚ_send_bytesr   Ú
recv_bytesÚ_recv_bytesÚpollÚ_pollr*   r%   r%   r&   r    C   s    


zQueue._after_forkTNc              	   C   sf   | j rtd| ›d�ƒ‚| j ||¡s(t‚| j�. | jd krB|  ¡  | j 	|¡ | j 
¡  W 5 Q R X d S ©NzQueue z
 is closed)r8   Ú
ValueErrorr   Úacquirer   r1   r5   Ú_start_threadr4   ÚappendÚnotify©r"   ÚobjÚblockÚtimeoutr%   r%   r&   ÚputP   s    
z	Queue.putc              	   C   sÄ   | j rtd| ›d�ƒ‚|rH|d krH| j� |  ¡ }W 5 Q R X | j ¡  nr|rXt ¡ | }| j ||¡sjt	‚zB|rŒ|t ¡  }|  
|¡s˜t	‚n|  
¡ s˜t	‚|  ¡ }| j ¡  W 5 | j ¡  X t |¡S r@   )r8   rA   r   r=   r   ÚreleaseÚtimeÚ	monotonicrB   r   r?   Ú_ForkingPicklerÚloads)r"   rH   rI   ÚresZdeadliner%   r%   r&   Úget\   s*    
z	Queue.getc                 C   s   | j | jj ¡  S r(   )r   r   Ú_semlockZ
_get_valuer*   r%   r%   r&   Úqsizev   s    zQueue.qsizec                 C   s
   |   ¡  S r(   ©r?   r*   r%   r%   r&   Úemptyz   s    zQueue.emptyc                 C   s   | j j ¡ S r(   )r   rR   Ú_is_zeror*   r%   r%   r&   Úfull}   s    z
Queue.fullc                 C   s
   |   d¡S ©NF)rQ   r*   r%   r%   r&   Ú
get_nowait€   s    zQueue.get_nowaitc                 C   s   |   |d¡S rX   )rJ   ©r"   rG   r%   r%   r&   Ú
put_nowaitƒ   s    zQueue.put_nowaitc                 C   s2   d| _ z| j ¡  W 5 | j}|r,d | _|ƒ  X d S )NT)r8   r9   r   Úclose)r"   r\   r%   r%   r&   r\   †   s    zQueue.closec                 C   s.   t dƒ | jstd | ¡ƒ‚| jr*|  ¡  d S )NzQueue.join_thread()zQueue {0!r} not closed)r
   r8   ÚAssertionErrorÚformatr6   r*   r%   r%   r&   Újoin_thread�   s    zQueue.join_threadc                 C   s6   t dƒ d| _z| j ¡  W n tk
r0   Y nX d S )NzQueue.cancel_join_thread()T)r
   r7   r6   ZcancelÚAttributeErrorr*   r%   r%   r&   Úcancel_join_thread–   s    zQueue.cancel_join_threadc              
   C   s°   t dƒ | j ¡  tjtj| j| j| j| j	| j
j| j| j| jfdd�| _d| j_t dƒ | j ¡  t dƒ | js�t| jtjt | j¡gdd�| _t| tj| j| jgd	d�| _d S )
NzQueue._start_thread()ZQueueFeederThread)ÚtargetÚargsÚnameTzdoing self._thread.start()z... done self._thread.start()éûÿÿÿ)Zexitpriorityé
   )r
   r4   Úclearr/   ZThreadr   Ú_feedr1   r;   r   r   r\   r   Ú_on_queue_feeder_errorr   r5   ZdaemonÚstartr7   r   Ú_finalize_joinÚweakrefÚrefr6   Ú_finalize_closer9   r*   r%   r%   r&   rC   ž   s<    
   þû
 ý 
ýzQueue._start_threadc                 C   s4   t dƒ | ƒ }|d k	r(| ¡  t dƒ nt dƒ d S )Nzjoining queue threadz... queue thread joinedz... queue thread already dead)r
   Újoin)ZtwrÚthreadr%   r%   r&   rk   ¾   s    
zQueue._finalize_joinc              	   C   s.   t dƒ |� |  t¡ | ¡  W 5 Q R X d S )Nztelling queue thread to quit)r
   rD   Ú	_sentinelrE   )ÚbufferÚnotemptyr%   r%   r&   rn   È   s    
zQueue._finalize_closec              
   C   sX  t dƒ |j}|j}	|j}
| j}t}tjdkr<|j}|j}nd }zš|ƒ  z| sT|
ƒ  W 5 |	ƒ  X zb|ƒ }||kr†t dƒ |ƒ  W W d S t 	|¡}|d kr¢||ƒ qb|ƒ  z||ƒ W 5 |ƒ  X qbW n t
k
rÖ   Y nX W q@ tk
�rP } zV|�rt|ddƒtjk�rW Y ¢6d S tƒ �r.td|ƒ W Y ¢d S | ¡  |||ƒ W 5 d }~X Y q@X q@d S )Nz$starting thread to feed data to piper   z%feeder thread got sentinel -- exitingÚerrnor   zerror in queue thread: %s)r
   rB   rK   ÚwaitÚpopleftrq   r   r   rN   ÚdumpsÚ
IndexErrorÚ	ExceptionÚgetattrrt   ZEPIPEr   r   )rr   rs   r:   Z	writelockr\   Zignore_epipeÚonerrorZ	queue_semZnacquireZnreleaseZnwaitZbpopleftÚsentinelZwacquireZwreleaserG   Úer%   r%   r&   rh   Ï   sN    







zQueue._feedc                 C   s   ddl }| ¡  dS )z˜
        Private API hook called when feeding data in the background thread
        raises an exception.  For overriding by concurrent.futures.
        r   N)Ú	tracebackÚ	print_exc)r}   rG   r~   r%   r%   r&   ri     s    zQueue._on_queue_feeder_error)r   )TN)TN)Ú__name__Ú
__module__Ú__qualname__r'   r+   r.   r    rJ   rQ   rS   rU   rW   rY   r[   r\   r_   ra   rC   Ústaticmethodrk   rn   rh   ri   r%   r%   r%   r&   r   "   s.   



 
	

=c                   @   s@   e Zd Zddd„Zdd„ Zdd„ Zdd
d„Zdd„ Zdd„ Zd	S )r   r   c                C   s*   t j| ||d� | d¡| _| ¡ | _d S )N)r$   r   )r   r'   Z	SemaphoreÚ_unfinished_tasksr0   Ú_condr!   r%   r%   r&   r'   #  s    zJoinableQueue.__init__c                 C   s   t  | ¡| j| jf S r(   )r   r+   r…   r„   r*   r%   r%   r&   r+   (  s    zJoinableQueue.__getstate__c                 C   s,   t  | |d d… ¡ |dd … \| _| _d S )Néþÿÿÿ)r   r.   r…   r„   r,   r%   r%   r&   r.   +  s    zJoinableQueue.__setstate__TNc              
   C   s‚   | j rtd| ›d�ƒ‚| j ||¡s(t‚| j�J | j�8 | jd krJ|  ¡  | j	 
|¡ | j ¡  | j ¡  W 5 Q R X W 5 Q R X d S r@   )r8   rA   r   rB   r   r1   r…   r5   rC   r4   rD   r„   rK   rE   rF   r%   r%   r&   rJ   /  s    

zJoinableQueue.putc              	   C   s@   | j �0 | j d¡stdƒ‚| jj ¡ r2| j  ¡  W 5 Q R X d S )NFz!task_done() called too many times)r…   r„   rB   rA   rR   rV   Z
notify_allr*   r%   r%   r&   Ú	task_done<  s
    zJoinableQueue.task_donec              	   C   s,   | j � | jj ¡ s| j  ¡  W 5 Q R X d S r(   )r…   r„   rR   rV   ru   r*   r%   r%   r&   ro   C  s    zJoinableQueue.join)r   )TN)	r€   r�   r‚   r'   r+   r.   rJ   r‡   ro   r%   r%   r%   r&   r   !  s   

c                   @   s<   e Zd Zdd„ Zdd„ Zdd„ Zdd„ Zd	d
„ Zdd„ ZdS )r   c                C   sH   t jdd�\| _| _| ¡ | _| jj| _tj	dkr:d | _
n
| ¡ | _
d S )NFr   r   )r   r   r   r   r   r   r>   r?   r   r   r   )r"   r$   r%   r%   r&   r'   N  s    


zSimpleQueue.__init__c                 C   s
   |   ¡  S r(   rT   r*   r%   r%   r&   rU   W  s    zSimpleQueue.emptyc                 C   s   t  | ¡ | j| j| j| jfS r(   )r	   r)   r   r   r   r   r*   r%   r%   r&   r+   Z  s    
zSimpleQueue.__getstate__c                 C   s"   |\| _ | _| _| _| j j| _d S r(   )r   r   r   r   r>   r?   r,   r%   r%   r&   r.   ^  s    zSimpleQueue.__setstate__c              	   C   s&   | j � | j ¡ }W 5 Q R X t |¡S r(   )r   r   r<   rN   rO   )r"   rP   r%   r%   r&   rQ   b  s    zSimpleQueue.getc              	   C   sD   t  |¡}| jd kr"| j |¡ n| j� | j |¡ W 5 Q R X d S r(   )rN   rw   r   r   r:   rZ   r%   r%   r&   rJ   h  s
    

zSimpleQueue.putN)	r€   r�   r‚   r'   rU   r+   r.   rQ   rJ   r%   r%   r%   r&   r   L  s   	)Ú__all__r   r   r/   r2   rL   rl   rt   Zqueuer   r   Z_multiprocessingÚ r   r	   Z	reductionZForkingPicklerrN   Úutilr
   r   r   r   r   Úobjectr   rq   r   r   r%   r%   r%   r&   Ú<module>
   s$   
 v
+