U
    ¼Ê¦iÛ¿  ã                   @   sÖ  d 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
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 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Ze
jdkrþedƒ‚dd„ ZG dd„ dejƒZG dd„ dej ƒZ!G dd„ dej"ej#ƒZ$G dd„ dej%ƒZ&G dd„ dƒZ'dd„ Z(G d d!„ d!e'ƒZ)G d"d#„ d#e)ƒZ*G d$d%„ d%e)ƒZ+G d&d'„ d'e'ƒZ,G d(d)„ d)e'ƒZ-G d*d+„ d+ej.ƒZ/eZ0e/Z1dS ),z2Selector event loop for Unix with signal handling.é    Né   )Úbase_events)Úbase_subprocess)Ú	constants)Ú
coroutines)Úevents)Ú
exceptions)Úfutures)Úselector_events)Útasks)Ú
transports)Úlogger)ÚSelectorEventLoopÚAbstractChildWatcherÚSafeChildWatcherÚFastChildWatcherÚMultiLoopChildWatcherÚThreadedChildWatcherÚDefaultEventLoopPolicyZwin32z+Signals are not really supported on Windowsc                 C   s   dS )zDummy signal handler.N© )ÚsignumÚframer   r   ú)/usr/lib/python3.8/asyncio/unix_events.pyÚ_sighandler_noop*   s    r   c                       sÊ   e Zd ZdZd)‡ fdd„	Z‡ fdd„Zdd„ Zd	d
„ Zdd„ Zdd„ Z	dd„ Z
d*dd„Zd+dd„Zd,dd„Zdd„ Zd-dddddœdd„Zd.ddddddœdd „Zd!d"„ Zd#d$„ Zd%d&„ Zd'd(„ Z‡  ZS )/Ú_UnixSelectorEventLoopzdUnix event loop.

    Adds signal handling and UNIX Domain Socket support to SelectorEventLoop.
    Nc                    s   t ƒ  |¡ i | _d S ©N)ÚsuperÚ__init__Ú_signal_handlers)ÚselfÚselector©Ú	__class__r   r   r   5   s    z_UnixSelectorEventLoop.__init__c                    sZ   t ƒ  ¡  t ¡ s.t| jƒD ]}|  |¡ qn(| jrVtjd| ›d�t	| d� | j 
¡  d S )NzClosing the loop z@ on interpreter shutdown stage, skipping signal handlers removal©Úsource)r   ÚcloseÚsysÚis_finalizingÚlistr   Úremove_signal_handlerÚwarningsÚwarnÚResourceWarningÚclear©r   Úsigr!   r   r   r%   9   s    
üz_UnixSelectorEventLoop.closec                 C   s   |D ]}|sq|   |¡ qd S r   )Ú_handle_signal)r   Údatar   r   r   r   Ú_process_self_dataG   s    z)_UnixSelectorEventLoop._process_self_datac                 G   sL  t  |¡st  |¡rtdƒ‚|  |¡ |  ¡  zt | j 	¡ ¡ W n2 t
tfk
rt } ztt|ƒƒ‚W 5 d}~X Y nX t ||| d¡}|| j|< zt |t¡ t |d¡ W nš tk
�rF } zz| j|= | j�szt d¡ W n4 t
tfk
�r } zt d|¡ W 5 d}~X Y nX |jtjk�r4td|› d�ƒ‚n‚ W 5 d}~X Y nX dS )zÃAdd a handler for a signal.  UNIX only.

        Raise ValueError if the signal number is invalid or uncatchable.
        Raise RuntimeError if there is a problem setting up the handler.
        z3coroutines cannot be used with add_signal_handler()NFéÿÿÿÿúset_wakeup_fd(-1) failed: %súsig ú cannot be caught)r   ZiscoroutineZiscoroutinefunctionÚ	TypeErrorÚ_check_signalZ_check_closedÚsignalÚset_wakeup_fdZ_csockÚfilenoÚ
ValueErrorÚOSErrorÚRuntimeErrorÚstrr   ZHandler   r   Úsiginterruptr   ÚinfoÚerrnoÚEINVAL)r   r/   ÚcallbackÚargsÚexcÚhandleZnexcr   r   r   Úadd_signal_handlerN   s2    
ÿ

z)_UnixSelectorEventLoop.add_signal_handlerc                 C   s8   | j  |¡}|dkrdS |jr*|  |¡ n
|  |¡ dS )z2Internal helper that is the actual signal handler.N)r   ÚgetZ
_cancelledr)   Z_add_callback_signalsafe)r   r/   rG   r   r   r   r0   {   s    z%_UnixSelectorEventLoop._handle_signalc              
   C   sæ   |   |¡ z| j|= W n tk
r,   Y dS X |tjkr@tj}ntj}zt ||¡ W nB tk
r˜ } z$|jtj	kr†t
d|› d�ƒ‚n‚ W 5 d}~X Y nX | jsâzt d¡ W n2 ttfk
rà } zt d|¡ W 5 d}~X Y nX dS )zwRemove a handler for a signal.  UNIX only.

        Return True if a signal handler was removed, False if not.
        Fr5   r6   Nr3   r4   T)r8   r   ÚKeyErrorr9   ÚSIGINTÚdefault_int_handlerÚSIG_DFLr=   rB   rC   r>   r:   r<   r   rA   )r   r/   ÚhandlerrF   r   r   r   r)   …   s(    

z,_UnixSelectorEventLoop.remove_signal_handlerc                 C   s6   t |tƒstd|›�ƒ‚|t ¡ kr2td|› �ƒ‚dS )zÁInternal helper to validate a signal.

        Raise ValueError if the signal number is invalid or uncatchable.
        Raise RuntimeError if there is a problem setting up the handler.
        zsig must be an int, not zinvalid signal number N)Ú
isinstanceÚintr7   r9   Úvalid_signalsr<   r.   r   r   r   r8   ¥   s    
z$_UnixSelectorEventLoop._check_signalc                 C   s   t | ||||ƒS r   )Ú_UnixReadPipeTransport©r   ÚpipeÚprotocolÚwaiterÚextrar   r   r   Ú_make_read_pipe_transport±   s    z0_UnixSelectorEventLoop._make_read_pipe_transportc                 C   s   t | ||||ƒS r   )Ú_UnixWritePipeTransportrS   r   r   r   Ú_make_write_pipe_transportµ   s    z1_UnixSelectorEventLoop._make_write_pipe_transportc	              
   Ë   s¼   t  ¡ �ª}
|
 ¡ stdƒ‚|  ¡ }t| |||||||f||dœ|	—Ž}|
 | ¡ | j|¡ z|I d H  W nD t	t
fk
r‚   ‚ Y n, tk
r¬   | ¡  | ¡ I d H  ‚ Y nX W 5 Q R X |S )NzRasyncio.get_child_watcher() is not activated, subprocess support is not installed.)rV   rW   )r   Úget_child_watcherÚ	is_activer>   Úcreate_futureÚ_UnixSubprocessTransportÚadd_child_handlerZget_pidÚ_child_watcher_callbackÚ
SystemExitÚKeyboardInterruptÚBaseExceptionr%   Z_wait)r   rU   rE   ÚshellÚstdinÚstdoutÚstderrÚbufsizerW   ÚkwargsÚwatcherrV   Útranspr   r   r   Ú_make_subprocess_transport¹   s8    

   ÿ þý
 ÿz1_UnixSelectorEventLoop._make_subprocess_transportc                 C   s   |   |j|¡ d S r   )Úcall_soon_threadsafeZ_process_exited)r   ÚpidÚ
returncoderk   r   r   r   r`   ×   s    z._UnixSelectorEventLoop._child_watcher_callback)ÚsslÚsockÚserver_hostnameÚssl_handshake_timeoutc          	      Ã   s   |d kst |tƒst‚|r,|d krLtdƒ‚n |d k	r<tdƒ‚|d k	rLtdƒ‚|d k	rº|d k	rdtdƒ‚t |¡}t tjtjd¡}z | 	d¡ |  
||¡I d H  W qú   | ¡  ‚ Y qúX n@|d krÊtdƒ‚|jtjksâ|jtjkrðtd|›�ƒ‚| 	d¡ | j|||||d	�I d H \}}||fS )
Nz/you have to pass server_hostname when using sslz+server_hostname is only meaningful with sslú1ssl_handshake_timeout is only meaningful with sslú3path and sock can not be specified at the same timer   Fzno path and sock were specifiedú.A UNIX Domain Stream Socket was expected, got )rs   )rO   r?   ÚAssertionErrorr<   ÚosÚfspathÚsocketÚAF_UNIXÚSOCK_STREAMÚsetblockingZsock_connectr%   ÚfamilyÚtypeZ_create_connection_transport)	r   Úprotocol_factoryÚpathrp   rq   rr   rs   Ú	transportrU   r   r   r   Úcreate_unix_connectionÚ   sT    ÿÿÿ



ÿÿ
   þz-_UnixSelectorEventLoop.create_unix_connectionéd   T)rq   Úbacklogrp   rs   Ústart_servingc             
   Ã   sÊ  t |tƒrtdƒ‚|d k	r&|s&tdƒ‚|d k	�rH|d k	r@tdƒ‚t |¡}t tjtj¡}|d dkrÊz t	 
t 	|¡j¡r„t |¡ W nB tk
rš   Y n0 tk
rÈ } zt d||¡ W 5 d }~X Y nX z| |¡ W nl tk
�r0 }	 z8| ¡  |	jtjk�rd|›d�}
ttj|
ƒd ‚n‚ W 5 d }	~	X Y n   | ¡  ‚ Y nX n<|d k�rZtd	ƒ‚|jtjk�sv|jtjk�r„td
|›�ƒ‚| d¡ t | |g||||¡}|�rÆ| ¡  tjd| d�I d H  |S )Nz*ssl argument must be an SSLContext or Nonert   ru   r   )r   ú z2Unable to check or remove stale UNIX socket %r: %rzAddress z is already in usez-path was not specified, and no sock specifiedrv   F)Úloop)rO   Úboolr7   r<   rx   ry   rz   r{   r|   ÚstatÚS_ISSOCKÚst_modeÚremoveÚFileNotFoundErrorr=   r   ÚerrorZbindr%   rB   Z
EADDRINUSEr~   r   r}   r   ZServerZ_start_servingr   Úsleep)r   r€   r�   rq   r…   rp   rs   r†   ÚerrrF   ÚmsgZserverr   r   r   Úcreate_unix_server  sn    
ÿ
ÿ
 ÿ

ÿ
ÿÿ
  ÿz)_UnixSelectorEventLoop.create_unix_serverc              
   Ã   sô   z
t j W n, tk
r6 } zt d¡‚W 5 d }~X Y nX z| ¡ }W n2 ttjfk
rv } zt d¡‚W 5 d }~X Y nX zt  |¡j	}W n, t
k
r´ } zt d¡‚W 5 d }~X Y nX |r¾|n|}	|	sÊdS |  ¡ }
|  |
d |||||	d¡ |
I d H S )Nzos.sendfile() is not availableznot a regular filer   )rx   ÚsendfileÚAttributeErrorr   ÚSendfileNotAvailableErrorr;   ÚioÚUnsupportedOperationÚfstatÚst_sizer=   r]   Ú_sock_sendfile_native_impl)r   rq   ÚfileÚoffsetÚcountrF   r;   r‘   ZfsizeÚ	blocksizeÚfutr   r   r   Ú_sock_sendfile_nativeJ  s2    
ÿ   ÿz,_UnixSelectorEventLoop._sock_sendfile_nativec	                 C   s,  |  ¡ }	|d k	r|  |¡ | ¡ r4|  |||¡ d S |rd|| }|dkrd|  |||¡ | |¡ d S zt |	|||¡}
W �nD ttfk
rÆ   |d kr¢|  	||¡ |  
|	| j||	||||||¡
 Y �nb tk
�rj } z†|d k	�r|jtjk�rt|ƒtk	�rtdtjƒ}||_|}|dk�rBt d¡}|  |||¡ | |¡ n|  |||¡ | |¡ W 5 d }~X Y n¾ ttfk
�r„   ‚ Y n¤ tk
�r¾ } z|  |||¡ | |¡ W 5 d }~X Y njX |
dk�rä|  |||¡ | |¡ nD||
7 }||
7 }|d k�r
|  	||¡ |  
|	| j||	||||||¡
 d S )Nr   zsocket is not connectedzos.sendfile call failed)r;   Úremove_writerÚ	cancelledÚ_sock_sendfile_update_fileposZ
set_resultrx   r”   ÚBlockingIOErrorÚInterruptedErrorÚ_sock_add_cancellation_callbackZ
add_writerr›   r=   rB   ZENOTCONNr   ÚConnectionErrorÚ	__cause__r   r–   Zset_exceptionra   rb   rc   )r   r    Zregistered_fdrq   r;   r�   rž   rŸ   Ú
total_sentÚfdZsentrF   Únew_excr‘   r   r   r   r›   a  s†    

     þ


ÿ
þ ÿ
ÿ

     þz1_UnixSelectorEventLoop._sock_sendfile_native_implc                 C   s   |dkrt  ||t j¡ d S ©Nr   )rx   ÚlseekÚSEEK_SET)r   r;   r�   rª   r   r   r   r¤   §  s    z4_UnixSelectorEventLoop._sock_sendfile_update_fileposc                    s   ‡ ‡fdd„}|  |¡ d S )Nc                    s&   |   ¡ r"ˆ ¡ }|dkr"ˆ  |¡ d S )Nr3   )r£   r;   r¢   )r    r«   ©r   rq   r   r   Úcb¬  s    zB_UnixSelectorEventLoop._sock_add_cancellation_callback.<locals>.cb)Zadd_done_callback)r   r    rq   r±   r   r°   r   r§   «  s    z6_UnixSelectorEventLoop._sock_add_cancellation_callback)N)NN)NN)N)N)N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r%   r2   rH   r0   r)   r8   rX   rZ   rl   r`   rƒ   r“   r¡   r›   r¤   r§   Ú__classcell__r   r   r!   r   r   /   sH   -
   ÿ
  ÿ
 þ
 ÿ ü. ÿ  üCFr   c                       sŠ   e Zd ZdZd‡ f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jfdd„Zddd„Zdd„ Zdd„ Z‡  ZS ) rR   i   Nc                    sÚ   t ƒ  |¡ || jd< || _|| _| ¡ | _|| _d| _d| _	t
 | j¡j}t |¡s„t |¡s„t |¡s„d | _d | _d | _tdƒ‚t
 | jd¡ | j | jj| ¡ | j | jj| j| j¡ |d k	rÖ| j tj|d ¡ d S )NrT   Fz)Pipe transport is for pipes/sockets only.)r   r   Ú_extraÚ_loopÚ_piper;   Ú_filenoÚ	_protocolÚ_closingÚ_pausedrx   r™   rŒ   rŠ   ÚS_ISFIFOr‹   ÚS_ISCHRr<   Úset_blockingÚ	call_soonÚconnection_madeÚ_add_readerÚ_read_readyr	   Ú_set_result_unless_cancelled)r   rˆ   rT   rU   rV   rW   Úmoder!   r   r   r   ¸  s:    


ÿþ ÿ
 ÿz_UnixReadPipeTransport.__init__c                 C   sÀ   | j jg}| jd kr | d¡ n| jr0| d¡ | d| j› �¡ t| jdd ƒ}| jd k	r�|d k	r�t 	|| jt
j¡}|r„| d¡ q°| d¡ n | jd k	r¦| d¡ n
| d¡ d d	 |¡¡S )
NÚclosedÚclosingúfd=Ú	_selectorÚpollingÚidleÚopenú<{}>ú )r"   r²   r¹   Úappendr¼   rº   Úgetattrr¸   r
   Ú_test_selector_eventÚ	selectorsZ
EVENT_READÚformatÚjoin)r   rA   r    rË   r   r   r   Ú__repr__Ö  s(    


  ÿ

z_UnixReadPipeTransport.__repr__c              
   C   sº   zt  | j| j¡}W nD ttfk
r,   Y nŠ tk
rX } z|  |d¡ W 5 d }~X Y n^X |rl| j 	|¡ nJ| j
 ¡ r‚t d| ¡ d| _| j
 | j¡ | j
 | jj¡ | j
 | jd ¡ d S )Nz"Fatal read error on pipe transportú%r was closed by peerT)rx   Úreadrº   Úmax_sizer¥   r¦   r=   Ú_fatal_errorr»   Zdata_receivedr¸   Ú	get_debugr   rA   r¼   Ú_remove_readerrÁ   Zeof_receivedÚ_call_connection_lost)r   r1   rF   r   r   r   rÄ   ë  s    
z"_UnixReadPipeTransport._read_readyc                 C   s>   | j s| jrd S d| _| j | j¡ | j ¡ r:t d| ¡ d S )NTz%r pauses reading)r¼   r½   r¸   rÜ   rº   rÛ   r   Údebug©r   r   r   r   Úpause_readingý  s    
z$_UnixReadPipeTransport.pause_readingc                 C   sB   | j s| jsd S d| _| j | j| j¡ | j ¡ r>t d| ¡ d S )NFz%r resumes reading)	r¼   r½   r¸   rÃ   rº   rÄ   rÛ   r   rÞ   rß   r   r   r   Úresume_reading  s    
z%_UnixReadPipeTransport.resume_readingc                 C   s
   || _ d S r   ©r»   ©r   rU   r   r   r   Úset_protocol  s    z#_UnixReadPipeTransport.set_protocolc                 C   s   | j S r   râ   rß   r   r   r   Úget_protocol  s    z#_UnixReadPipeTransport.get_protocolc                 C   s   | j S r   ©r¼   rß   r   r   r   Ú
is_closing  s    z!_UnixReadPipeTransport.is_closingc                 C   s   | j s|  d ¡ d S r   )r¼   Ú_closerß   r   r   r   r%     s    z_UnixReadPipeTransport.closec                 C   s,   | j d k	r(|d| ›�t| d� | j  ¡  d S ©Nzunclosed transport r#   ©r¹   r,   r%   ©r   Ú_warnr   r   r   Ú__del__  s    
z_UnixReadPipeTransport.__del__úFatal error on pipe transportc                 C   sZ   t |tƒr4|jtjkr4| j ¡ rLtjd| |dd� n| j ||| | j	dœ¡ |  
|¡ d S ©Nz%r: %sT©Úexc_info)ÚmessageÚ	exceptionr‚   rU   )rO   r=   rB   ZEIOr¸   rÛ   r   rÞ   Úcall_exception_handlerr»   rè   ©r   rF   rò   r   r   r   rÚ     s    
üz#_UnixReadPipeTransport._fatal_errorc                 C   s(   d| _ | j | j¡ | j | j|¡ d S ©NT)r¼   r¸   rÜ   rº   rÁ   rÝ   ©r   rF   r   r   r   rè   -  s    z_UnixReadPipeTransport._closec                 C   s4   z| j |¡ W 5 | j  ¡  d | _ d | _d | _X d S r   ©r¹   r%   r»   r¸   Zconnection_lostr÷   r   r   r   rÝ   2  s    
z,_UnixReadPipeTransport._call_connection_lost)NN)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   rR   ´  s   
rR   c                       s¨   e Zd Zd%‡ f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ejfdd„Zdd„ Zd&dd „Zd'd!d"„Zd#d$„ Z‡  ZS )(rY   Nc           
         sþ   t ƒ  ||¡ || jd< || _| ¡ | _|| _tƒ | _d| _	d| _
t | j¡j}t |¡}t |¡}t |¡}	|s”|s”|	s”d | _d | _d | _tdƒ‚t | jd¡ | j | jj| ¡ |	sÈ|ràtj d¡sà| j | jj| j| j¡ |d k	rú| j tj|d ¡ d S )NrT   r   Fz?Pipe transport is only for pipes, sockets and character devicesZaix)r   r   r·   r¹   r;   rº   r»   Ú	bytearrayÚ_bufferÚ
_conn_lostr¼   rx   r™   rŒ   rŠ   r¿   r¾   r‹   r<   rÀ   r¸   rÁ   rÂ   r&   ÚplatformÚ
startswithrÃ   rÄ   r	   rÅ   )
r   rˆ   rT   rU   rV   rW   rÆ   Zis_charZis_fifoZ	is_socketr!   r   r   r   ?  s:    




 ÿ
 ÿz _UnixWritePipeTransport.__init__c                 C   sØ   | j jg}| jd kr | d¡ n| jr0| d¡ | d| j› �¡ t| jdd ƒ}| jd k	r¨|d k	r¨t 	|| jt
j¡}|r„| d¡ n
| d¡ |  ¡ }| d|› �¡ n | jd k	r¾| d¡ n
| d¡ d	 d
 |¡¡S )NrÇ   rÈ   rÉ   rÊ   rË   rÌ   zbufsize=rÍ   rÎ   rÏ   )r"   r²   r¹   rÐ   r¼   rº   rÑ   r¸   r
   rÒ   rÓ   ZEVENT_WRITEÚget_write_buffer_sizerÔ   rÕ   )r   rA   r    rË   rh   r   r   r   rÖ   d  s,    


  ÿ


z _UnixWritePipeTransport.__repr__c                 C   s
   t | jƒS r   )Úlenrú   rß   r   r   r   rþ   |  s    z-_UnixWritePipeTransport.get_write_buffer_sizec                 C   s6   | j  ¡ rt d| ¡ | jr*|  tƒ ¡ n|  ¡  d S )Nr×   )r¸   rÛ   r   rA   rú   rè   ÚBrokenPipeErrorrß   r   r   r   rÄ     s
    
z#_UnixWritePipeTransport._read_readyc              
   C   sR  t |tttfƒstt|ƒƒ‚t |tƒr.t|ƒ}|s6d S | jsB| jrj| jtj	krXt
 d¡ |  jd7  _d S | j�s8zt | j|¡}W nt ttfk
r    d}Y nZ ttfk
r¸   ‚ Y nB tk
rø } z$|  jd7  _|  |d¡ W Y ¢d S d }~X Y nX |t|ƒk�rd S |dk�r&t|ƒ|d … }| j | j| j¡ |  j|7  _|  ¡  d S )Nz=pipe closed by peer or os.write(pipe, data) raised exception.r   r   ú#Fatal write error on pipe transport)rO   Úbytesrù   Ú
memoryviewrw   Úreprrû   r¼   r   Z!LOG_THRESHOLD_FOR_CONNLOST_WRITESr   Úwarningrú   rx   Úwriterº   r¥   r¦   ra   rb   rc   rÚ   rÿ   r¸   Z_add_writerÚ_write_readyZ_maybe_pause_protocol)r   r1   ÚnrF   r   r   r   r  ˆ  s8    


z_UnixWritePipeTransport.writec              
   C   s  | j stdƒ‚zt | j| j ¡}W n‚ ttfk
r:   Y nÒ ttfk
rR   ‚ Y nº t	k
r¤ } z6| j  
¡  |  jd7  _| j | j¡ |  |d¡ W 5 d }~X Y nhX |t| j ƒkrö| j  
¡  | j | j¡ |  ¡  | jrò| j | j¡ |  d ¡ d S |dk�r| j d |…= d S )NzData should not be emptyr   r  r   )rú   rw   rx   r  rº   r¥   r¦   ra   rb   rc   r-   rû   r¸   Ú_remove_writerrÚ   rÿ   Z_maybe_resume_protocolr¼   rÜ   rÝ   )r   r  rF   r   r   r   r  «  s,    



z$_UnixWritePipeTransport._write_readyc                 C   s   dS rö   r   rß   r   r   r   Úcan_write_eofÇ  s    z%_UnixWritePipeTransport.can_write_eofc                 C   sB   | j r
d S | jst‚d| _ | js>| j | j¡ | j | jd ¡ d S rö   )	r¼   r¹   rw   rú   r¸   rÜ   rº   rÁ   rÝ   rß   r   r   r   Ú	write_eofÊ  s    
z!_UnixWritePipeTransport.write_eofc                 C   s
   || _ d S r   râ   rã   r   r   r   rä   Ó  s    z$_UnixWritePipeTransport.set_protocolc                 C   s   | j S r   râ   rß   r   r   r   rå   Ö  s    z$_UnixWritePipeTransport.get_protocolc                 C   s   | j S r   ræ   rß   r   r   r   rç   Ù  s    z"_UnixWritePipeTransport.is_closingc                 C   s   | j d k	r| js|  ¡  d S r   )r¹   r¼   r  rß   r   r   r   r%   Ü  s    z_UnixWritePipeTransport.closec                 C   s,   | j d k	r(|d| ›�t| d� | j  ¡  d S ré   rê   rë   r   r   r   rí   á  s    
z_UnixWritePipeTransport.__del__c                 C   s   |   d ¡ d S r   )rè   rß   r   r   r   Úabortæ  s    z_UnixWritePipeTransport.abortrî   c                 C   sN   t |tƒr(| j ¡ r@tjd| |dd� n| j ||| | jdœ¡ |  |¡ d S rï   )	rO   r=   r¸   rÛ   r   rÞ   rô   r»   rè   rõ   r   r   r   rÚ   é  s    

üz$_UnixWritePipeTransport._fatal_errorc                 C   sF   d| _ | jr| j | j¡ | j ¡  | j | j¡ | j | j|¡ d S rö   )	r¼   rú   r¸   r	  rº   r-   rÜ   rÁ   rÝ   r÷   r   r   r   rè   ÷  s    
z_UnixWritePipeTransport._closec                 C   s4   z| j |¡ W 5 | j  ¡  d | _ d | _d | _X d S r   rø   r÷   r   r   r   rÝ   ÿ  s    
z-_UnixWritePipeTransport._call_connection_lost)NN)rî   )N)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   rY   <  s"   %	#	

rY   c                   @   s   e Zd Zdd„ ZdS )r^   c           	   	   K   sŠ   d }|t jkrt ¡ \}}zPt j|f||||d|dœ|—Ž| _|d k	rh| ¡  t| ¡ d|d�| j_	d }W 5 |d k	r„| ¡  | ¡  X d S )NF)rd   re   rf   rg   Zuniversal_newlinesrh   Úwb)Ú	buffering)
Ú
subprocessÚPIPErz   Z
socketpairr%   ÚPopenÚ_procrÍ   Údetachre   )	r   rE   rd   re   rf   rg   rh   ri   Zstdin_wr   r   r   Ú_start  s.    
ÿ    þþz_UnixSubprocessTransport._startN)r²   r³   r´   r  r   r   r   r   r^   	  s   r^   c                   @   sH   e 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S )r   aH  Abstract base class for monitoring child processes.

    Objects derived from this class monitor a collection of subprocesses and
    report their termination or interruption by a signal.

    New callbacks are registered with .add_child_handler(). Starting a new
    process must be done within a 'with' block to allow the watcher to suspend
    its activity until the new process if fully registered (this is needed to
    prevent a race condition in some implementations).

    Example:
        with watcher:
            proc = subprocess.Popen("sleep 1")
            watcher.add_child_handler(proc.pid, callback)

    Notes:
        Implementations of this class must be thread-safe.

        Since child watcher objects may catch the SIGCHLD signal and call
        waitpid(-1), there should be only one active object per process.
    c                 G   s
   t ƒ ‚dS )a  Register a new child handler.

        Arrange for callback(pid, returncode, *args) to be called when
        process 'pid' terminates. Specifying another callback for the same
        process replaces the previous handler.

        Note: callback() must be thread-safe.
        N©ÚNotImplementedError©r   rn   rD   rE   r   r   r   r_   9  s    	z&AbstractChildWatcher.add_child_handlerc                 C   s
   t ƒ ‚dS )z Removes the handler for process 'pid'.

        The function returns True if the handler was successfully removed,
        False if there was nothing to remove.Nr  ©r   rn   r   r   r   Úremove_child_handlerD  s    z)AbstractChildWatcher.remove_child_handlerc                 C   s
   t ƒ ‚dS )zÔAttach the watcher to an event loop.

        If the watcher was previously attached to an event loop, then it is
        first detached before attaching to the new loop.

        Note: loop may be None.
        Nr  ©r   rˆ   r   r   r   Úattach_loopL  s    z AbstractChildWatcher.attach_loopc                 C   s
   t ƒ ‚dS )zlClose the watcher.

        This must be called to make sure that any underlying resource is freed.
        Nr  rß   r   r   r   r%   V  s    zAbstractChildWatcher.closec                 C   s
   t ƒ ‚dS )zºReturn ``True`` if the watcher is active and is used by the event loop.

        Return True if the watcher is installed and ready to handle process exit
        notifications.

        Nr  rß   r   r   r   r\   ]  s    zAbstractChildWatcher.is_activec                 C   s
   t ƒ ‚dS )zdEnter the watcher's context and allow starting new processes

        This function must return selfNr  rß   r   r   r   Ú	__enter__f  s    zAbstractChildWatcher.__enter__c                 C   s
   t ƒ ‚dS )zExit the watcher's contextNr  ©r   ÚaÚbÚcr   r   r   Ú__exit__l  s    zAbstractChildWatcher.__exit__N)r²   r³   r´   rµ   r_   r  r  r%   r\   r  r!  r   r   r   r   r   "  s   
	r   c                 C   s2   t  | ¡rt  | ¡ S t  | ¡r*t  | ¡S | S d S r   )rx   ÚWIFSIGNALEDÚWTERMSIGÚ	WIFEXITEDÚWEXITSTATUS)Ústatusr   r   r   Ú_compute_returncodeq  s
    


r'  c                   @   sD   e 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S )ÚBaseChildWatcherc                 C   s   d | _ i | _d S r   )r¸   Ú
_callbacksrß   r   r   r   r   �  s    zBaseChildWatcher.__init__c                 C   s   |   d ¡ d S r   )r  rß   r   r   r   r%   …  s    zBaseChildWatcher.closec                 C   s   | j d k	o| j  ¡ S r   )r¸   Z
is_runningrß   r   r   r   r\   ˆ  s    zBaseChildWatcher.is_activec                 C   s
   t ƒ ‚d S r   r  )r   Úexpected_pidr   r   r   Ú_do_waitpid‹  s    zBaseChildWatcher._do_waitpidc                 C   s
   t ƒ ‚d S r   r  rß   r   r   r   Ú_do_waitpid_allŽ  s    z BaseChildWatcher._do_waitpid_allc                 C   s~   |d kst |tjƒst‚| jd k	r<|d kr<| jr<t dt¡ | jd k	rT| j 	t
j¡ || _|d k	rz| t
j| j¡ |  ¡  d S )NzCA loop is being detached from a child watcher with pending handlers)rO   r   ZAbstractEventLooprw   r¸   r)  r*   r+   ÚRuntimeWarningr)   r9   ÚSIGCHLDrH   Ú	_sig_chldr,  r  r   r   r   r  ‘  s    ý
zBaseChildWatcher.attach_loopc              
   C   s^   z|   ¡  W nL ttfk
r&   ‚ Y n4 tk
rX } z| j d|dœ¡ W 5 d }~X Y nX d S )Nú$Unknown exception in SIGCHLD handler)rò   ró   )r,  ra   rb   rc   r¸   rô   r÷   r   r   r   r/  ¥  s    þzBaseChildWatcher._sig_chldN)
r²   r³   r´   r   r%   r\   r+  r,  r  r/  r   r   r   r   r(    s   r(  c                       sP   e Zd ZdZ‡ fdd„Zdd„ Zdd„ Zdd	„ Zd
d„ Zdd„ Z	dd„ Z
‡  ZS )r   ad  'Safe' child watcher implementation.

    This implementation avoids disrupting other code spawning processes by
    polling explicitly each process in the SIGCHLD handler instead of calling
    os.waitpid(-1).

    This is a safe solution but it has a significant overhead when handling a
    big number of children (O(n) each time SIGCHLD is raised)
    c                    s   | j  ¡  tƒ  ¡  d S r   )r)  r-   r   r%   rß   r!   r   r   r%   ¿  s    
zSafeChildWatcher.closec                 C   s   | S r   r   rß   r   r   r   r  Ã  s    zSafeChildWatcher.__enter__c                 C   s   d S r   r   r  r   r   r   r!  Æ  s    zSafeChildWatcher.__exit__c                 G   s   ||f| j |< |  |¡ d S r   )r)  r+  r  r   r   r   r_   É  s    z"SafeChildWatcher.add_child_handlerc                 C   s*   z| j |= W dS  tk
r$   Y dS X d S ©NTF©r)  rJ   r  r   r   r   r  Ï  s
    z%SafeChildWatcher.remove_child_handlerc                 C   s   t | jƒD ]}|  |¡ q
d S r   ©r(   r)  r+  r  r   r   r   r,  Ö  s    z SafeChildWatcher._do_waitpid_allc                 C   sÐ   |dkst ‚zt |tj¡\}}W n( tk
rJ   |}d}t d|¡ Y n.X |dkrXd S t|ƒ}| j 	¡ rxt 
d||¡ z| j |¡\}}W n. tk
rº   | j 	¡ r¶tjd|dd� Y nX |||f|žŽ  d S )Nr   éÿ   ú8Unknown child process pid %d, will report returncode 255ú$process %s exited with returncode %sú'Child watcher got an unexpected pid: %rTrð   )rw   rx   ÚwaitpidÚWNOHANGÚChildProcessErrorr   r  r'  r¸   rÛ   rÞ   r)  ÚpoprJ   )r   r*  rn   r&  ro   rD   rE   r   r   r   r+  Û  s6    þ

 ÿ
 ÿzSafeChildWatcher._do_waitpid)r²   r³   r´   rµ   r%   r  r!  r_   r  r,  r+  r¶   r   r   r!   r   r   ´  s   
r   c                       sT   e Zd ZdZ‡ fdd„Z‡ fdd„Zdd„ Zdd	„ Zd
d„ Zdd„ Z	dd„ Z
‡  ZS )r   aW  'Fast' child watcher implementation.

    This implementation reaps every terminated processes by calling
    os.waitpid(-1) directly, possibly breaking other code spawning processes
    and waiting for their termination.

    There is no noticeable overhead when handling a big number of children
    (O(1) each time a child terminates).
    c                    s$   t ƒ  ¡  t ¡ | _i | _d| _d S r­   )r   r   Ú	threadingZLockÚ_lockÚ_zombiesÚ_forksrß   r!   r   r   r     s    

zFastChildWatcher.__init__c                    s"   | j  ¡  | j ¡  tƒ  ¡  d S r   )r)  r-   r>  r   r%   rß   r!   r   r   r%     s    

zFastChildWatcher.closec              
   C   s0   | j �  |  jd7  _| W  5 Q R £ S Q R X d S )Nr   )r=  r?  rß   r   r   r   r    s    zFastChildWatcher.__enter__c              	   C   s^   | j �B |  jd8  _| js"| js0W 5 Q R £ d S t| jƒ}| j ¡  W 5 Q R X t d|¡ d S )Nr   z5Caught subprocesses termination from unknown pids: %s)r=  r?  r>  r?   r-   r   r  )r   r  r  r   Zcollateral_victimsr   r   r   r!    s    
þzFastChildWatcher.__exit__c              	   G   st   | j stdƒ‚| j�F z| j |¡}W n. tk
rT   ||f| j|< Y W 5 Q R £ d S X W 5 Q R X |||f|žŽ  d S )NzMust use the context manager)r?  rw   r=  r>  r;  rJ   r)  )r   rn   rD   rE   ro   r   r   r   r_   '  s    z"FastChildWatcher.add_child_handlerc                 C   s*   z| j |= W dS  tk
r$   Y dS X d S r1  r2  r  r   r   r   r  5  s
    z%FastChildWatcher.remove_child_handlerc              	   C   sþ   zt  dt j¡\}}W n tk
r,   Y d S X |dkr:d S t|ƒ}| j�‚ z| j |¡\}}W nN tk
r¬   | j	r¤|| j
|< | j ¡ r–t d||¡ Y W 5 Q R £ q d }Y nX | j ¡ rÆt d||¡ W 5 Q R X |d krèt d||¡ q |||f|žŽ  q d S )Nr3   r   z,unknown process %s exited with returncode %sr6  z8Caught subprocess termination from unknown pid: %d -> %d)rx   r8  r9  r:  r'  r=  r)  r;  rJ   r?  r>  r¸   rÛ   r   rÞ   r  )r   rn   r&  ro   rD   rE   r   r   r   r,  <  s@    

 þ

 ÿ þz FastChildWatcher._do_waitpid_all)r²   r³   r´   rµ   r   r%   r  r!  r_   r  r,  r¶   r   r   r!   r   r   þ  s   	r   c                   @   sh   e 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S )r   a~  A watcher that doesn't require running loop in the main thread.

    This implementation registers a SIGCHLD signal handler on
    instantiation (which may conflict with other code that
    install own handler for this signal).

    The solution is safe but it has a significant overhead when
    handling a big number of processes (*O(n)* each time a
    SIGCHLD is received).
    c                 C   s   i | _ d | _d S r   )r)  Ú_saved_sighandlerrß   r   r   r   r   z  s    zMultiLoopChildWatcher.__init__c                 C   s
   | j d k	S r   )r@  rß   r   r   r   r\   ~  s    zMultiLoopChildWatcher.is_activec                 C   sT   | j  ¡  | jd krd S t tj¡}|| jkr:t d¡ nt tj| j¡ d | _d S )Nz+SIGCHLD handler was changed by outside code)	r)  r-   r@  r9   Ú	getsignalr.  r/  r   r  )r   rN   r   r   r   r%   �  s    


zMultiLoopChildWatcher.closec                 C   s   | S r   r   rß   r   r   r   r  �  s    zMultiLoopChildWatcher.__enter__c                 C   s   d S r   r   ©r   Úexc_typeZexc_valZexc_tbr   r   r   r!  �  s    zMultiLoopChildWatcher.__exit__c                 G   s&   t  ¡ }|||f| j|< |  |¡ d S r   )r   Úget_running_loopr)  r+  )r   rn   rD   rE   rˆ   r   r   r   r_   “  s    z'MultiLoopChildWatcher.add_child_handlerc                 C   s*   z| j |= W dS  tk
r$   Y dS X d S r1  r2  r  r   r   r   r  š  s
    z*MultiLoopChildWatcher.remove_child_handlerc                 C   sN   | j d k	rd S t tj| j¡| _ | j d kr<t d¡ tj| _ t tjd¡ d S )NzaPrevious SIGCHLD handler was set by non-Python code, restore to default handler on watcher close.F)r@  r9   r.  r/  r   r  rM   r@   r  r   r   r   r  ¡  s    


z!MultiLoopChildWatcher.attach_loopc                 C   s   t | jƒD ]}|  |¡ q
d S r   r3  r  r   r   r   r,  ²  s    z%MultiLoopChildWatcher._do_waitpid_allc           	      C   sî   |dkst ‚zt |tj¡\}}W n, tk
rN   |}d}t d|¡ d}Y nX |dkr\d S t|ƒ}d}z| j 	|¡\}}}W n$ t
k
r¢   tjd|dd� Y nHX | ¡ r¼t d||¡ n.|rÖ| ¡ rÖt d	||¡ |j|||f|žŽ  d S )
Nr   r4  r5  FTr7  rð   ú%Loop %r that handles pid %r is closedr6  )rw   rx   r8  r9  r:  r   r  r'  r)  r;  rJ   Ú	is_closedrÛ   rÞ   rm   )	r   r*  rn   r&  ro   Z	debug_logrˆ   rD   rE   r   r   r   r+  ¶  s<    þ
 ÿ ÿz!MultiLoopChildWatcher._do_waitpidc              	   C   sL   z|   ¡  W n: ttfk
r&   ‚ Y n" tk
rF   tjddd� Y nX d S )Nr0  Trð   )r,  ra   rb   rc   r   r  )r   r   r   r   r   r   r/  Û  s    zMultiLoopChildWatcher._sig_chldN)r²   r³   r´   rµ   r   r\   r%   r  r!  r_   r  r  r,  r+  r/  r   r   r   r   r   g  s   %r   c                   @   sn   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
d„ Zdd„ Z	e
jfdd„Zdd„ Zdd„ Zdd„ Zdd„ ZdS )r   aA  Threaded child watcher implementation.

    The watcher uses a thread per process
    for waiting for the process finish.

    It doesn't require subscription on POSIX signal
    but a thread creation is not free.

    The watcher has O(1) complexity, its performance doesn't depend
    on amount of spawn processes.
    c                 C   s   t  d¡| _i | _d S r­   )Ú	itertoolsrž   Ú_pid_counterÚ_threadsrß   r   r   r   r   ñ  s    zThreadedChildWatcher.__init__c                 C   s   dS rö   r   rß   r   r   r   r\   õ  s    zThreadedChildWatcher.is_activec                 C   s   |   ¡  d S r   )Ú_join_threadsrß   r   r   r   r%   ø  s    zThreadedChildWatcher.closec                 C   s.   dd„ t | j ¡ ƒD ƒ}|D ]}| ¡  qdS )z%Internal: Join all non-daemon threadsc                 S   s   g | ]}|  ¡ r|js|‘qS r   )Úis_aliveÚdaemon©Ú.0Úthreadr   r   r   Ú
<listcomp>ý  s     ÿz6ThreadedChildWatcher._join_threads.<locals>.<listcomp>N)r(   rI  ÚvaluesrÕ   )r   ÚthreadsrO  r   r   r   rJ  û  s    z"ThreadedChildWatcher._join_threadsc                 C   s   | S r   r   rß   r   r   r   r    s    zThreadedChildWatcher.__enter__c                 C   s   d S r   r   rB  r   r   r   r!    s    zThreadedChildWatcher.__exit__c                 C   s6   dd„ t | j ¡ ƒD ƒ}|r2|| j› d�t| d� d S )Nc                 S   s   g | ]}|  ¡ r|‘qS r   )rK  rM  r   r   r   rP  	  s    ÿz0ThreadedChildWatcher.__del__.<locals>.<listcomp>z0 has registered but not finished child processesr#   )r(   rI  rQ  r"   r,   )r   rì   rR  r   r   r   rí     s    þzThreadedChildWatcher.__del__c                 G   sF   t  ¡ }tj| jdt| jƒ› �||||fdd�}|| j|< | ¡  d S )Nzwaitpid-T)ÚtargetÚnamerE   rL  )	r   rD  r<  ZThreadr+  ÚnextrH  rI  Ústart)r   rn   rD   rE   rˆ   rO  r   r   r   r_     s    
ý
z&ThreadedChildWatcher.add_child_handlerc                 C   s   dS rö   r   r  r   r   r   r    s    z)ThreadedChildWatcher.remove_child_handlerc                 C   s   d S r   r   r  r   r   r   r    s    z ThreadedChildWatcher.attach_loopc                 C   s¤   |dkst ‚zt |d¡\}}W n( tk
rH   |}d}t d|¡ Y n X t|ƒ}| ¡ rht d||¡ | 	¡ r€t d||¡ n|j
|||f|žŽ  | j |¡ d S )Nr   r4  r5  r6  rE  )rw   rx   r8  r:  r   r  r'  rÛ   rÞ   rF  rm   rI  r;  )r   rˆ   r*  rD   rE   rn   r&  ro   r   r   r   r+  "  s(    þ
 ÿz ThreadedChildWatcher._do_waitpidN)r²   r³   r´   rµ   r   r\   r%   rJ  r  r!  r*   r+   rí   r_   r  r  r+  r   r   r   r   r   ä  s   	r   c                       sH   e Zd ZdZeZ‡ fdd„Zdd„ Z‡ fdd„Zdd	„ Z	d
d„ Z
‡  ZS )Ú_UnixDefaultEventLoopPolicyz:UNIX event loop policy with a watcher for child processes.c                    s   t ƒ  ¡  d | _d S r   )r   r   Ú_watcherrß   r!   r   r   r   A  s    
z$_UnixDefaultEventLoopPolicy.__init__c              	   C   sH   t j�8 | jd kr:tƒ | _tt ¡ tjƒr:| j | j	j
¡ W 5 Q R X d S r   )r   r=  rX  r   rO   r<  Úcurrent_threadÚ_MainThreadr  Ú_localr¸   rß   r   r   r   Ú_init_watcherE  s    
ÿz)_UnixDefaultEventLoopPolicy._init_watcherc                    s6   t ƒ  |¡ | jdk	r2tt ¡ tjƒr2| j |¡ dS )zÑSet the event loop.

        As a side effect, if a child watcher was set before, then calling
        .set_event_loop() from the main thread will call .attach_loop(loop) on
        the child watcher.
        N)r   Úset_event_looprX  rO   r<  rY  rZ  r  r  r!   r   r   r]  M  s
    
ÿz*_UnixDefaultEventLoopPolicy.set_event_loopc                 C   s   | j dkr|  ¡  | j S )z~Get the watcher for child processes.

        If not yet set, a ThreadedChildWatcher object is automatically created.
        N)rX  r\  rß   r   r   r   r[   [  s    
z-_UnixDefaultEventLoopPolicy.get_child_watcherc                 C   s4   |dkst |tƒst‚| jdk	r*| j ¡  || _dS )z$Set the watcher for child processes.N)rO   r   rw   rX  r%   )r   rj   r   r   r   Úset_child_watchere  s    

z-_UnixDefaultEventLoopPolicy.set_child_watcher)r²   r³   r´   rµ   r   Z_loop_factoryr   r\  r]  r[   r^  r¶   r   r   r!   r   rW  =  s   
rW  )2rµ   rB   r—   rG  rx   rÓ   r9   rz   rŠ   r  r&   r<  r*   Ú r   r   r   r   r   r   r	   r
   r   r   Úlogr   Ú__all__rü   ÚImportErrorr   ZBaseSelectorEventLoopr   ZReadTransportrR   Z_FlowControlMixinZWriteTransportrY   ZBaseSubprocessTransportr^   r   r'  r(  r   r   r   r   ZBaseDefaultEventLoopPolicyrW  r   r   r   r   r   r   Ú<module>   s`   	
    	ÿ NO5Ji}Y3