
Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­
<!DOCTYPE html>
<html>
3

  \”‘  ã               @   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
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dlmZ dddddgZejdkrüedƒ‚dd„ Zy
ejZW n ek
�r,   dd„ ZY nX G dd„ dejƒZ e!edƒ�rVdd„ Z"nddl#Z#d d„ Z"G d!d"„ d"ej$ƒZ%G d#d$„ d$ej&ej'ƒZ(e!ed%ƒ�r¢ej)Z*nddl#Z#d&d'„ Z*G d(d)„ d)ej+ƒZ,G d*d„ dƒZ-G d+d,„ d,e-ƒZ.G d-d„ de.ƒZ/G d.d„ de.ƒZ0G d/d0„ d0ej1ƒZ2e Z3e2Z4dS )1z2Selector event loop for Unix with signal handling.é    Né   )Úbase_events)Úbase_subprocess)Úcompat)Ú	constants)Ú
coroutines)Úevents)Úfutures)Úselector_events)Ú	selectors)Ú
transports)Ú	coroutine)ÚloggerÚSelectorEventLoopÚAbstractChildWatcherÚSafeChildWatcherÚFastChildWatcherÚDefaultEventLoopPolicyZwin32z+Signals are not really supported on Windowsc             C   s   dS )zDummy signal handler.N© )ÚsignumÚframer   r   ú+/usr/lib64/python3.6/asyncio/unix_events.pyÚ_sighandler_noop%   s    r   c             C   s   | S )Nr   )Úpathr   r   r   Ú<lambda>.   s    r   c                   s¶   e Zd ZdZd"‡ fdd„	Z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ed%dd„ƒZdd„ Zeddddœdd„ƒZed&ddddœd d!„ƒZ‡  ZS )'Ú_UnixSelectorEventLoopzdUnix event loop.

    Adds signal handling and UNIX Domain Socket support to SelectorEventLoop.
    Nc                s   t ƒ j|ƒ i | _d S )N)ÚsuperÚ__init__Ú_signal_handlers)ÚselfÚselector)Ú	__class__r   r   r   7   s    z_UnixSelectorEventLoop.__init__c             C   s   t jƒ S )N)ÚsocketZ
socketpair)r   r   r   r   Ú_socketpair;   s    z"_UnixSelectorEventLoop._socketpairc                s^   t ƒ jƒ  tjƒ s2xFt| jƒD ]}| j|ƒ qW n(| jrZtjd| ›d�t	| d� | j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   Úsig)r!   r   r   r%   >   s    
z_UnixSelectorEventLoop.closec             C   s"   x|D ]}|sq| j |ƒ qW d S )N)Ú_handle_signal)r   Údatar   r   r   r   Ú_process_self_dataL   s    
z)_UnixSelectorEventLoop._process_self_datac          +   G   sH  t j|ƒst j|ƒrtdƒ‚| j|ƒ | jƒ  ytj| jj	ƒ ƒ W n2 t
tfk
rt } ztt|ƒƒ‚W Y dd}~X nX tj||| ƒ}|| j|< ytj|tƒ tj|dƒ W n˜ tk
�rB } zz| j|= | j�sytjdƒ W n4 t
tfk
�r } ztjd|ƒ W Y dd}~X nX |jtjk�r0tdj|ƒƒ‚n‚ W Y dd}~X 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()NFr   zset_wakeup_fd(-1) failed: %sz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Úformat)r   r.   ÚcallbackÚargsÚexcÚhandleZnexcr   r   r   Úadd_signal_handlerS   s0    



z)_UnixSelectorEventLoop.add_signal_handlerc             C   s8   | j j|ƒ}|dkrdS |jr*| j|ƒ n
| j|ƒ dS )z2Internal helper that is the actual signal handler.N)r   ÚgetZ
_cancelledr)   Z_add_callback_signalsafe)r   r.   rD   r   r   r   r/   €   s    z%_UnixSelectorEventLoop._handle_signalc          &   C   sâ   | j |ƒ y| j|= W n tk
r*   dS X |tjkr>tj}ntj}ytj||ƒ W n@ tk
r” } z$|jtj	kr‚t
dj|ƒƒ‚n‚ W Y dd}~X nX | jsÞytjdƒ W n2 ttfk
rÜ } ztjd|ƒ W Y dd}~X nX dS )zwRemove a handler for a signal.  UNIX only.

        Return True if a signal handler was removed, False if not.
        Fzsig {} cannot be caughtNr   zset_wakeup_fd(-1) failed: %sTr2   )r4   r   ÚKeyErrorr5   ÚSIGINTÚdefault_int_handlerÚSIG_DFLr9   r>   r?   r:   r@   r6   r8   r   r=   )r   r.   ZhandlerrC   r   r   r   r)   Š   s(    

z,_UnixSelectorEventLoop.remove_signal_handlerc             C   sH   t |tƒstdj|ƒƒ‚d|  ko,tjk n  sDtdj|tjƒƒ‚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 {!r}r   zsig {} out of range(1, {})N)Ú
isinstanceÚintr3   r@   r5   ÚNSIGr8   )r   r.   r   r   r   r4   ª   s
    
z$_UnixSelectorEventLoop._check_signalc             C   s   t | ||||ƒS )N)Ú_UnixReadPipeTransport)r   ÚpipeÚprotocolÚwaiterÚextrar   r   r   Ú_make_read_pipe_transport·   s    z0_UnixSelectorEventLoop._make_read_pipe_transportc             C   s   t | ||||ƒS )N)Ú_UnixWritePipeTransport)r   rO   rP   rQ   rR   r   r   r   Ú_make_write_pipe_transport»   s    z1_UnixSelectorEventLoop._make_write_pipe_transportc	             k   s´   t jƒ �¢}
| jƒ }t| |||||||f||dœ|	—Ž}|
j|jƒ | j|ƒ y|E d H  W n& tk
r~ } z
|}W Y d d }~X nX d }|d k	r¦|jƒ  |j	ƒ E d H  |‚W d Q R X |S )N)rQ   rR   )
r   Úget_child_watcherZcreate_futureÚ_UnixSubprocessTransportÚadd_child_handlerZget_pidÚ_child_watcher_callbackÚ	Exceptionr%   Z_wait)r   rP   rB   ÚshellÚstdinÚstdoutÚstderrÚbufsizerR   ÚkwargsÚwatcherrQ   ÚtransprC   Úerrr   r   r   Ú_make_subprocess_transport¿   s$    




z1_UnixSelectorEventLoop._make_subprocess_transportc             C   s   | j |j|ƒ d S )N)Zcall_soon_threadsafeZ_process_exited)r   ÚpidÚ
returncoderb   r   r   r   rY   Ý   s    z._UnixSelectorEventLoop._child_watcher_callback)ÚsslÚsockÚserver_hostnamec            c   s  |d kst |tƒst‚|r,|d kr<tdƒ‚n|d k	r<tdƒ‚|d k	r |d k	rTtdƒ‚tjtjtjdƒ}y |jdƒ | j||ƒE d H  W qâ   |j	ƒ  ‚ Y qâX nB|d kr°tdƒ‚|j
tjksÊtj|jƒ rØtdj|ƒƒ‚|jdƒ | j||||ƒE d H \}}||fS )Nz/you have to pass server_hostname when using sslz+server_hostname is only meaningful with sslz3path and sock can not be specified at the same timer   Fzno path and sock were specifiedz2A UNIX Domain Stream Socket was expected, got {!r})rK   r;   ÚAssertionErrorr8   r"   ÚAF_UNIXÚSOCK_STREAMÚsetblockingZsock_connectr%   Úfamilyr   Ú_is_stream_socketÚtyper@   Z_create_connection_transport)r   Úprotocol_factoryr   rg   rh   ri   Ú	transportrP   r   r   r   Úcreate_unix_connectionà   s:    


z-_UnixSelectorEventLoop.create_unix_connectionéd   )rh   Úbacklogrg   c      
   !   C   s¤  t |tƒrtdƒ‚|d k	�r0|d k	r,tdƒ‚t|ƒ}tjtjtjƒ}|d d
kr´y tj	t
j|ƒjƒrnt
j|ƒ W nB tk
r„   Y n0 tk
r² } ztjd||ƒ W Y d d }~X nX y|j|ƒ W nj tk
�r } z8|jƒ  |jtjk�rdj|ƒ}ttj|ƒd ‚n‚ W Y d d }~X n   |jƒ  ‚ Y nX n>|d k�rBtdƒ‚|jtjk�s`tj|jƒ �rntdj|ƒƒ‚tj| |gƒ}	|j|ƒ |jd	ƒ | j||||	ƒ |	S )Nz*ssl argument must be an SSLContext or Nonez3path and sock can not be specified at the same timer   ú z2Unable to check or remove stale UNIX socket %r: %rzAddress {!r} is already in usez-path was not specified, and no sock specifiedz2A UNIX Domain Stream Socket was expected, got {!r}F)r   rv   )rK   Úboolr3   r8   Ú_fspathr"   rk   rl   ÚstatÚS_ISSOCKÚosÚst_modeÚremoveÚFileNotFoundErrorr9   r   ÚerrorZbindr%   r>   Z
EADDRINUSEr@   rn   r   ro   rp   ZServerZlistenrm   Z_start_serving)
r   rq   r   rh   ru   rg   rc   rC   ÚmsgZserverr   r   r   Úcreate_unix_server  sP    

 




z)_UnixSelectorEventLoop.create_unix_server)N)NN)NN)N)N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r#   r%   r1   rE   r/   r)   r4   rS   rU   r   rd   rY   rs   r�   Ú__classcell__r   r   )r!   r   r   1   s,   -
  
 
%r   Úset_blockingc             C   s   t j| dƒ d S )NF)r{   r‡   )Úfdr   r   r   Ú_set_nonblockingB  s    r‰   c             C   s,   t j | t jƒ}|tjB }t j | t j|ƒ d S )N)ÚfcntlZF_GETFLr{   Ú
O_NONBLOCKZF_SETFL)rˆ   Úflagsr   r   r   r‰   G  s    
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rhdd„ Zd!dd„Zdd„ Zdd„ Z‡  ZS )"rN   é   i   Nc                sÐ   t ƒ j|ƒ || jd< || _|| _|jƒ | _|| _d| _t	j
| jƒj}tj|ƒpbtj|ƒpbtj|ƒs~d | _d | _d | _tdƒ‚t| jƒ | jj| jj| ƒ | jj| jj| j| jƒ |d k	rÌ| jjtj|d ƒ d S )NrO   Fz)Pipe transport is for pipes/sockets only.)r   r   Ú_extraÚ_loopÚ_piper7   Ú_filenoÚ	_protocolÚ_closingr{   Úfstatr|   ry   ÚS_ISFIFOrz   ÚS_ISCHRr8   r‰   Ú	call_soonÚconnection_madeÚ_add_readerÚ_read_readyr	   Ú_set_result_unless_cancelled)r   ÚlooprO   rP   rQ   rR   Úmode)r!   r   r   r   Q  s,    






z_UnixReadPipeTransport.__init__c             C   s¼   | j jg}| jd kr |jdƒ n| jr0|jdƒ |jd| j ƒ t| jdd ƒ}| jd k	rŽ|d k	rŽtj	|| jt
jƒ}|r‚|jdƒ q®|jdƒ n | jd k	r¤|jdƒ n
|jdƒ dd	j|ƒ S )
NÚclosedÚclosingzfd=%sÚ	_selectorÚpollingÚidleÚopenz<%s>ú )r!   r‚   r�   Úappendr“   r‘   Úgetattrr�   r
   Ú_test_selector_eventr   Z
EVENT_READÚjoin)r   r=   r    r¡   r   r   r   Ú__repr__n  s$    




z_UnixReadPipeTransport.__repr__c             C   sº   yt j| j| jƒ}W nD ttfk
r,   Y nŠ tk
rX } z| j|dƒ W Y d d }~X n^X |rl| jj	|ƒ nJ| j
jƒ r‚tjd| ƒ d| _| j
j| jƒ | j
j| jjƒ | j
j| jd ƒ d S )Nz"Fatal read error on pipe transportz%r was closed by peerT)r{   Úreadr‘   Úmax_sizeÚBlockingIOErrorÚInterruptedErrorr9   Ú_fatal_errorr’   Zdata_receivedr�   Ú	get_debugr   r=   r“   Ú_remove_readerr—   Zeof_receivedÚ_call_connection_lost)r   r0   rC   r   r   r   rš   „  s    
z"_UnixReadPipeTransport._read_readyc             C   s   | j j| jƒ d S )N)r�   r°   r‘   )r   r   r   r   Úpause_reading–  s    z$_UnixReadPipeTransport.pause_readingc             C   s   | j j| j| jƒ d S )N)r�   r™   r‘   rš   )r   r   r   r   Úresume_reading™  s    z%_UnixReadPipeTransport.resume_readingc             C   s
   || _ d S )N)r’   )r   rP   r   r   r   Úset_protocolœ  s    z#_UnixReadPipeTransport.set_protocolc             C   s   | j S )N)r’   )r   r   r   r   Úget_protocolŸ  s    z#_UnixReadPipeTransport.get_protocolc             C   s   | j S )N)r“   )r   r   r   r   Ú
is_closing¢  s    z!_UnixReadPipeTransport.is_closingc             C   s   | j s| jd ƒ d S )N)r“   Ú_close)r   r   r   r   r%   ¥  s    z_UnixReadPipeTransport.closec             C   s,   | j d k	r(tjd|  t| d� | j jƒ  d S )Nzunclosed transport %r)r$   )r�   r*   r+   r,   r%   )r   r   r   r   Ú__del__­  s    
z_UnixReadPipeTransport.__del__úFatal error on pipe transportc             C   sZ   t |tƒr4|jtjkr4| jjƒ rLtjd| |dd� n| jj||| | j	dœƒ | j
|ƒ d S )Nz%r: %sT)Úexc_info)ÚmessageÚ	exceptionrr   rP   )rK   r9   r>   ZEIOr�   r¯   r   ÚdebugÚcall_exception_handlerr’   r·   )r   rC   r»   r   r   r   r®   ³  s    
z#_UnixReadPipeTransport._fatal_errorc             C   s(   d| _ | jj| jƒ | jj| j|ƒ d S )NT)r“   r�   r°   r‘   r—   r±   )r   rC   r   r   r   r·   Á  s    z_UnixReadPipeTransport._closec             C   s4   z| j j|ƒ W d | jjƒ  d | _d | _ d | _X d S )N)r’   Úconnection_lostr�   r%   r�   )r   rC   r   r   r   r±   Æ  s    
z,_UnixReadPipeTransport._call_connection_losti   )NN)r¹   )r‚   rƒ   r„   r«   r   r©   rš   r²   r³   r´   rµ   r¶   r%   r   ÚPY34r¸   r®   r·   r±   r†   r   r   )r!   r   rN   M  s   
rN   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r|dd„ Zdd„ Zd&dd „Zd'd!d"„Zd#d$„ Z‡  ZS )(rT   Nc       
         sü   t ƒ j||ƒ || jd< || _|jƒ | _|| _tƒ | _d| _	d| _
tj| jƒj}tj|ƒ}tj|ƒ}tj|ƒ}	|px|px|	s”d | _d | _d | _tdƒ‚t| jƒ | jj| jj| ƒ |	sÆ|rÞtjjdƒ rÞ| jj| jj| j| jƒ |d k	rø| jjtj|d ƒ d S )NrO   r   Fz?Pipe transport is only for pipes, sockets and character devicesÚaix)r   r   rŽ   r�   r7   r‘   r’   Ú	bytearrayÚ_bufferÚ
_conn_lostr“   r{   r”   r|   ry   r–   r•   rz   r8   r‰   r�   r—   r˜   r&   ÚplatformÚ
startswithr™   rš   r	   r›   )
r   rœ   rO   rP   rQ   rR   r�   Zis_charZis_fifoZ	is_socket)r!   r   r   r   Ó  s2    






z _UnixWritePipeTransport.__init__c             C   sÒ   | j jg}| jd kr |jdƒ n| jr0|jdƒ |jd| j ƒ t| jdd ƒ}| jd k	r¤|d k	r¤tj	|| jt
jƒ}|r‚|jdƒ n
|jdƒ | jƒ }|jd| ƒ n | jd k	rº|jdƒ n
|jdƒ d	d
j|ƒ S )Nrž   rŸ   zfd=%sr    r¡   r¢   z
bufsize=%sr£   z<%s>r¤   )r!   r‚   r�   r¥   r“   r‘   r¦   r�   r
   r§   r   ZEVENT_WRITEÚget_write_buffer_sizer¨   )r   r=   r    r¡   r_   r   r   r   r©   ø  s(    





z _UnixWritePipeTransport.__repr__c             C   s
   t | jƒS )N)ÚlenrÃ   )r   r   r   r   rÇ     s    z-_UnixWritePipeTransport.get_write_buffer_sizec             C   s6   | j jƒ rtjd| ƒ | jr*| jtƒ ƒ n| jƒ  d S )Nz%r was closed by peer)r�   r¯   r   r=   rÃ   r·   ÚBrokenPipeError)r   r   r   r   rš     s
    
z#_UnixWritePipeTransport._read_readyc             C   s0  t |tttfƒstt|ƒƒ‚t |tƒr.t|ƒ}|s6d S | jsB| jrj| jtj	krXt
jdƒ |  jd7  _d S | j�sytj| j|ƒ}W nT ttfk
r    d}Y n: tk
rØ } z|  jd7  _| j|dƒ d S d }~X nX |t|ƒkrêd S |dk�rt|ƒ|d … }| jj| j| jƒ |  j|7  _| jƒ  d S )Nz=pipe closed by peer or os.write(pipe, data) raised exception.r   r   z#Fatal write error on pipe transport)rK   ÚbytesrÂ   Ú
memoryviewrj   ÚreprrÄ   r“   r   Z!LOG_THRESHOLD_FOR_CONNLOST_WRITESr   ÚwarningrÃ   r{   Úwriter‘   r¬   r­   rZ   r®   rÈ   r�   Z_add_writerÚ_write_readyZ_maybe_pause_protocol)r   r0   ÚnrC   r   r   r   rÎ     s4    


z_UnixWritePipeTransport.writec             C   sö   | j stdƒ‚ytj| j| j ƒ}W nj ttfk
r:   Y n¸ tk
rŒ } z8| j jƒ  |  j	d7  _	| j
j| jƒ | j|dƒ W Y d d }~X nfX |t| j ƒkrÞ| j jƒ  | j
j| jƒ | jƒ  | jrÚ| j
j| jƒ | jd ƒ d S |dkrò| j d |…= d S )NzData should not be emptyr   z#Fatal write error on pipe transportr   )rÃ   rj   r{   rÎ   r‘   r¬   r­   rZ   r-   rÄ   r�   Ú_remove_writerr®   rÈ   Z_maybe_resume_protocolr“   r°   r±   )r   rÐ   rC   r   r   r   rÏ   >  s(    


z$_UnixWritePipeTransport._write_readyc             C   s   dS )NTr   )r   r   r   r   Úcan_write_eofX  s    z%_UnixWritePipeTransport.can_write_eofc             C   sB   | j r
d S | jst‚d| _ | js>| jj| jƒ | jj| jd ƒ d S )NT)	r“   r�   rj   rÃ   r�   r°   r‘   r—   r±   )r   r   r   r   Ú	write_eof[  s    
z!_UnixWritePipeTransport.write_eofc             C   s
   || _ d S )N)r’   )r   rP   r   r   r   r´   d  s    z$_UnixWritePipeTransport.set_protocolc             C   s   | j S )N)r’   )r   r   r   r   rµ   g  s    z$_UnixWritePipeTransport.get_protocolc             C   s   | j S )N)r“   )r   r   r   r   r¶   j  s    z"_UnixWritePipeTransport.is_closingc             C   s   | j d k	r| j r| jƒ  d S )N)r�   r“   rÓ   )r   r   r   r   r%   m  s    z_UnixWritePipeTransport.closec             C   s,   | j d k	r(tjd|  t| d� | j jƒ  d S )Nzunclosed transport %r)r$   )r�   r*   r+   r,   r%   )r   r   r   r   r¸   v  s    
z_UnixWritePipeTransport.__del__c             C   s   | j d ƒ d S )N)r·   )r   r   r   r   Úabort|  s    z_UnixWritePipeTransport.abortúFatal error on pipe transportc             C   sP   t |tjƒr*| jjƒ rBtjd| |dd� n| jj||| | jdœƒ | j	|ƒ d S )Nz%r: %sT)rº   )r»   r¼   rr   rP   )
rK   r   Z_FATAL_ERROR_IGNOREr�   r¯   r   r½   r¾   r’   r·   )r   rC   r»   r   r   r   r®     s    
z$_UnixWritePipeTransport._fatal_errorc             C   sF   d| _ | jr| jj| jƒ | jjƒ  | jj| jƒ | jj| j|ƒ d S )NT)	r“   rÃ   r�   rÑ   r‘   r-   r°   r—   r±   )r   rC   r   r   r   r·   �  s    
z_UnixWritePipeTransport._closec             C   s4   z| j j|ƒ W d | jjƒ  d | _d | _ d | _X d S )N)r’   r¿   r�   r%   r�   )r   rC   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   rT   Ð  s$   %	!	

rT   Úset_inheritablec             C   sN   t tddƒ}tj| tjƒ}|s4tj| tj||B ƒ ntj| tj|| @ ƒ d S )NZ
FD_CLOEXECr   )r¦   rŠ   ZF_GETFDZF_SETFD)rˆ   ZinheritableZcloexec_flagÚoldr   r   r   Ú_set_inheritable¥  s
    rØ   c               @   s   e Zd Zdd„ ZdS )rW   c       	   	   K   sv   d }|t jkr*| jjƒ \}}t|jƒ dƒ t j|f||||d|dœ|—Ž| _|d k	rr|jƒ  t	|j
ƒ d|d�| j_d S )NF)r[   r\   r]   r^   Zuniversal_newlinesr_   Úwb)Ú	buffering)Ú
subprocessÚPIPEr�   r#   rØ   r7   ÚPopenÚ_procr%   r£   Údetachr\   )	r   rB   r[   r\   r]   r^   r_   r`   Zstdin_wr   r   r   Ú_start±  s    
z_UnixSubprocessTransport._startN)r‚   rƒ   r„   rà   r   r   r   r   rW   ¯  s   rW   c               @   s@   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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   re   rA   rB   r   r   r   rX   ß  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.N)rá   )r   re   r   r   r   Úremove_child_handlerê  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.
        N)rá   )r   rœ   r   r   r   Úattach_loopò  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.
        N)rá   )r   r   r   r   r%   ü  s    zAbstractChildWatcher.closec             C   s
   t ƒ ‚dS )zdEnter the watcher's context and allow starting new processes

        This function must return selfN)rá   )r   r   r   r   Ú	__enter__  s    zAbstractChildWatcher.__enter__c             C   s
   t ƒ ‚dS )zExit the watcher's contextN)rá   )r   ÚaÚbÚcr   r   r   Ú__exit__	  s    zAbstractChildWatcher.__exit__N)
r‚   rƒ   r„   r…   rX   râ   rã   r%   rä   rè   r   r   r   r   r   È  s   
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 )N)r�   Ú
_callbacks)r   r   r   r   r     s    zBaseChildWatcher.__init__c             C   s   | j d ƒ d S )N)rã   )r   r   r   r   r%     s    zBaseChildWatcher.closec             C   s
   t ƒ ‚d S )N)rá   )r   Úexpected_pidr   r   r   Ú_do_waitpid  s    zBaseChildWatcher._do_waitpidc             C   s
   t ƒ ‚d S )N)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jdtƒ | jd k	rT| jj	t
jƒ || _|d k	rz|jt
j| jƒ | jƒ  d S )NzCA loop is being detached from a child watcher with pending handlers)rK   r   ZAbstractEventLooprj   r�   rê   r*   r+   ÚRuntimeWarningr)   r5   ÚSIGCHLDrE   Ú	_sig_chldrí   )r   rœ   r   r   r   rã     s    
zBaseChildWatcher.attach_loopc             C   sF   y| j ƒ  W n4 tk
r@ } z| jjd|dœƒ W Y d d }~X nX d S )Nz$Unknown exception in SIGCHLD handler)r»   r¼   )rí   rZ   r�   r¾   )r   rC   r   r   r   rð   1  s    zBaseChildWatcher._sig_chldc             C   s2   t j|ƒrt j|ƒ S t j|ƒr*t j|ƒS |S d S )N)r{   ÚWIFSIGNALEDÚWTERMSIGÚ	WIFEXITEDÚWEXITSTATUS)r   Ústatusr   r   r   Ú_compute_returncode=  s
    


z$BaseChildWatcher._compute_returncodeN)
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 jƒ  tƒ jƒ  d S )N)rê   r-   r   r%   )r   )r!   r   r   r%   V  s    
zSafeChildWatcher.closec             C   s   | S )Nr   )r   r   r   r   rä   Z  s    zSafeChildWatcher.__enter__c             C   s   d S )Nr   )r   rå   ræ   rç   r   r   r   rè   ]  s    zSafeChildWatcher.__exit__c             G   s.   | j d krtdƒ‚||f| j|< | j|ƒ d S )NzICannot add child handler, the child watcher does not have a loop attached)r�   r:   rê   rì   )r   re   rA   rB   r   r   r   rX   `  s
    
z"SafeChildWatcher.add_child_handlerc             C   s&   y| j |= dS  tk
r    dS X d S )NTF)rê   rG   )r   re   r   r   r   râ   k  s
    z%SafeChildWatcher.remove_child_handlerc             C   s"   xt | jƒD ]}| j|ƒ qW d S )N)r(   rê   rì   )r   re   r   r   r   rí   r  s    z SafeChildWatcher._do_waitpid_allc             C   sÒ   |dkst ‚ytj|tjƒ\}}W n( tk
rJ   |}d}tjd|ƒ Y n0X |dkrXd S | j|ƒ}| jj	ƒ rztj
d||ƒ y| jj|ƒ\}}W n. tk
r¼   | jj	ƒ r¸tjd|dd� Y nX |||f|žŽ  d S )Nr   éÿ   z8Unknown child process pid %d, will report returncode 255z$process %s exited with returncode %sz'Child watcher got an unexpected pid: %rT)rº   )rj   r{   ÚwaitpidÚWNOHANGÚChildProcessErrorr   rÍ   rö   r�   r¯   r½   rê   ÚpoprG   )r   rë   re   rõ   rf   rA   rB   r   r   r   rì   w  s,    


zSafeChildWatcher._do_waitpid)r‚   rƒ   r„   r…   r%   rä   rè   rX   râ   rí   rì   r†   r   r   )r!   r   r   K  s   	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 ƒ jƒ  tjƒ | _i | _d| _d S )Nr   )r   r   Ú	threadingZLockÚ_lockÚ_zombiesÚ_forks)r   )r!   r   r   r   ¤  s    

zFastChildWatcher.__init__c                s"   | j jƒ  | jjƒ  tƒ jƒ  d S )N)rê   r-   rþ   r   r%   )r   )r!   r   r   r%   ª  s    

zFastChildWatcher.closec          
   C   s$   | j � |  jd7  _| S Q R X d S )Nr   )rý   rÿ   )r   r   r   r   rä   ¯  s    zFastChildWatcher.__enter__c          
   C   sV   | j �: |  jd8  _| js$| j r(d S t| jƒ}| jjƒ  W d Q R X tj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   sz   | j stdƒ‚| jd kr tdƒ‚| j�: y| jj|ƒ}W n" tk
rZ   ||f| j|< d S X W d Q R X |||f|žŽ  d S )NzMust use the context managerzICannot add child handler, the child watcher does not have a loop attached)	rÿ   rj   r�   r:   rý   rþ   rû   rG   rê   )r   re   rA   rB   rf   r   r   r   rX   Ã  s    
z"FastChildWatcher.add_child_handlerc             C   s&   y| j |= dS  tk
r    dS X d S )NTF)rê   rG   )r   re   r   r   r   râ   Ö  s
    z%FastChildWatcher.remove_child_handlerc             C   sö   xðyt jdt jƒ\}}W n tk
r,   d S X |dkr:d S | j|ƒ}| j�v y| jj|ƒ\}}W nB tk
r¢   | j	rš|| j
|< | jjƒ r˜tjd||ƒ wd }Y nX | jjƒ r¼tjd||ƒ W d Q R X |d krÞtjd||ƒ q|||f|žŽ  qW d S )Nr   r   z,unknown process %s exited with returncode %sz$process %s exited with returncode %sz8Caught subprocess termination from unknown pid: %d -> %dr2   )r{   rø   rù   rú   rö   rý   rê   rû   rG   rÿ   rþ   r�   r¯   r   r½   rÍ   )r   re   rõ   rf   rA   rB   r   r   r   rí   Ý  s6    





z FastChildWatcher._do_waitpid_all)r‚   rƒ   r„   r…   r   r%   rä   rè   rX   râ   rí   r†   r   r   )r!   r   r   š  s   	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 ƒ jƒ  d | _d S )N)r   r   Ú_watcher)r   )r!   r   r   r     s    
z$_UnixDefaultEventLoopPolicy.__init__c          
   C   sH   t j�8 | jd kr:tƒ | _ttjƒ tjƒr:| jj| j	j
ƒ W d Q R X d S )N)r   rý   r  r   rK   rü   Úcurrent_threadÚ_MainThreadrã   Ú_localr�   )r   r   r   r   Ú_init_watcher  s    
z)_UnixDefaultEventLoopPolicy._init_watcherc                s6   t ƒ j|ƒ | jdk	r2ttjƒ tjƒr2| j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_loopr  rK   rü   r  r  rã   )r   rœ   )r!   r   r   r    s    
z*_UnixDefaultEventLoopPolicy.set_event_loopc             C   s   | j dkr| jƒ  | j S )zzGet the watcher for child processes.

        If not yet set, a SafeChildWatcher object is automatically created.
        N)r  r  )r   r   r   r   rV   &  s    
z-_UnixDefaultEventLoopPolicy.get_child_watcherc             C   s4   |dkst |tƒst‚| jdk	r*| jjƒ  || _dS )z$Set the watcher for child processes.N)rK   r   rj   r  r%   )r   ra   r   r   r   Úset_child_watcher0  s    

z-_UnixDefaultEventLoopPolicy.set_child_watcher)r‚   rƒ   r„   r…   r   Z_loop_factoryr   r  r  rV   r  r†   r   r   )r!   r   r     s   
r   )5r…   r>   r{   r5   r"   ry   rÛ   r&   rü   r*   Ú r   r   r   r   r   r   r	   r
   r   r   r   Úlogr   Ú__all__rÅ   ÚImportErrorr   Úfspathrx   ÚAttributeErrorZBaseSelectorEventLoopr   Úhasattrr‰   rŠ   ZReadTransportrN   Z_FlowControlMixinZWriteTransportrT   rÖ   rØ   ZBaseSubprocessTransportrW   r   ré   r   r   ZBaseDefaultEventLoopPolicyr   r   r   r   r   r   r   Ú<module>   sn   

  
  O
F=On2