o
    y.e}                     @   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zd dlZW n e	y-   dZY nw ej
dkrId dlmZ d dlZd dlZd dlZd dlZejdZdd ZG dd deZG dd	 d	eZd
d ZG dd deZede G dd deZede G dd deZerede dS dS )    Nwin32streamc                 C   s.   t | }|r| S t| rt| S dS )a   True if the stream or pstream specified by 'name' needs periodic probes
    to verify connectivity.  For [p]streams which need probes, it can take a
    long time to notice the connection was dropped.  Returns False if probes
    aren't needed, and None if 'name' is invalidN)Stream_find_methodneeds_probesPassiveStreamis_valid_name)namecls r   ,/usr/lib/python3/dist-packages/ovs/stream.pystream_or_pstream_needs_probes'   s   


r   c                   @   sh  e Zd ZdZdZdZdZdZdZdZ	i Z
dZdZdZdZdZdZdZdZedd Zed	d
 Zedd Zd@ddZdZed? Zedd ZeefddZedd ZedAd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)d0d1 Z*d2d3 Z+d4d5 Z,d6d7 Z-d8d9 Z.ed:d; Z/ed<d= Z0ed>d? Z1dS )Br   zQBidirectional byte stream.  Unix domain sockets, tcp and ssl
    are implemented.r         NFc                 C   s   |t j| d < d S )N:)r   _SOCKET_METHODS)methodr
   r   r   r   register_methodQ      zStream.register_methodc                 C   s*   t j D ]\}}| |r|  S qd S N)r   r   items
startswith)r	   r   r
   r   r   r   r   U   s
   
zStream._find_methodc                 C   s   t t| S )zReturns True if 'name' is a stream name in the form "TYPE:ARGS" and
        TYPE is a supported stream type ("unix:", "tcp:" and "ssl:"),
        otherwise False.)boolr   r   r	   r   r   r   r   \   s   zStream.is_valid_namec                 C   s   || _ || _tjdkrH|d ur@|| _|ddd }tjtj	j
|}t|| _t | _t | j_t | _t | j_ntjddd| _|| _|tjkrUtj| _n|dkr^tj| _ntj| _d| _d S )Nr   r   r   F)bManualResetbInitialStater   )socketpipesysplatform_serversplitovsutilabs_file_namedirsRUNDIRwinutilsget_pipe_name	_pipename
pywintypes
OVERLAPPED_readget_new_eventhEvent_write_weventr	   errnoEAGAINr   _Stream__S_CONNECTINGstate_Stream__S_CONNECTED_Stream__S_DISCONNECTEDerror)selfr   r	   statusr   	is_serversuffixr   r   r   __init__c   s,   






zStream.__init__   c                 C   s   t j| S r   )r"   socket_utilcheck_connection_completion)sockr   r   r   r?      s   z"Stream.check_connection_completionc                 C   s  t | }|stjdfS | ddd }| drtjtj	j
|}tjdkrt|}t|dkr7tjdfS z	t|d  W n
   tjdf Y S z!t|}z	t|tj W n tjyj   tjdf Y W S w W n5 tjy } z(|jtjjkrdt _d	|d| tjtjd
dfW  Y d}~S tjdfW  Y d}~S d}~ww d	|d| d	|d
dfS |||\}}|r|dfS | |}	|	tjks|	tj!krtj}
d	}	n	|	d	krd	}
n|	}
|	||| |
fS )aq  Attempts to connect a stream to a remote peer.  'name' is a
        connection name in the form "TYPE:ARGS", where TYPE is an active stream
        class's name and ARGS are stream class-specific.  The supported TYPEs
        include "unix", "tcp", and "ssl".

        Returns (error, stream): on success 'error' is 0 and 'stream' is the
        new Stream, on failure 'error' is a positive errno value and 'stream'
        is None.

        Never returns errno.EAGAIN or errno.EINPROGRESS.  Instead, returns 0
        and a new Stream.  The connect() method can be used to check for
        successful connection completion.Nr   r   zunix:r      rTr   F)r   r:   )"r   r   r1   EAFNOSUPPORTr!   r   r"   r#   r$   r%   r&   r   r   r'   r(   lenENOENTopenclosecreate_fileset_pipe_mode	win32pipePIPE_READMODE_BYTEr*   r7   winerrorERROR_PIPE_BUSYretry_connectr2   	win32fileINVALID_HANDLE_VALUE_openr?   EINPROGRESS)r	   dscpr
   r;   pipenamenpipeer7   r@   errr9   r   r   r   rF      s\   






	
zStream.openc                 C   s   t d)Nz)This method must be overrided by subclass)NotImplementedError)r;   rS   r   r   r   rQ      s   zStream._openc                 C   s   | \}}|skd}|dur|dkrt j | }	 | }tjdkr)|tjkr)tj}|tjkr/n0|dur>t j |kr>tj	}n!|
  t j }|| || |durZ|| |  q|jdurk|tjkskJ |ru|ru|  d}||fS )a  Blocks until a Stream completes its connection attempt, either
        succeeding or failing, but no more than 'timeout' milliseconds.
        (error, stream) should be the tuple returned by Stream.open().
        Negative value of 'timeout' means infinite waiting.
        Returns a tuple of the same form.

        Typical usage:
        error, stream = Stream.open_block(Stream.open("unix:/tmp/socket"))Nr   Tr   )r"   timevalmsecconnectr   r   r1   WSAEWOULDBLOCKr2   	ETIMEDOUTrunpollerPollerrun_waitconnect_waittimer_wait_untilblockr   rR   rG   )error_streamtimeoutr7   r   deadliner_   r   r   r   
open_block   s8   





zStream.open_blockc                 C   sx   | j d ur
| j   | jd ur:| jrt| j t| j t| jt	j
 t| jjt	j
 t| jjt	j
 d S d S r   )r   rG   r   r    rJ   FlushFileBuffersDisconnectNamedPiper'   close_handlevlogwarnr,   r.   r/   r8   r   r   r   rG      s   


zStream.closec              
   C   s   | j d ur| | j }|tjksJ n=tjdkrP| jrNzt| j	| _
d| _d}W n& tjyM } z|jtjjkr=tj}nd| _tj}W Y d }~nd }~ww d}|dkrZtj| _d S |tjkrhtj| _|| _d S d S )Nr   Fr   )r   r?   r1   rR   r   r   rN   r'   rH   r)   r   _retry_connectr*   r7   rL   rM   r2   rE   r   r5   r4   r6   r8   retvalrV   r   r   r   __scs_connecting
  s.   

	

zStream.__scs_connectingc                 C   sL   | j tjkr
|   | j tjkrtjS | j tjkrdS | j tjks#J | jS )zTries to complete the connection on this stream.  If the connection
        is complete, returns 0 if the connection was successful or a positive
        errno value if it failed.  If the connection is still in progress,
        returns errno.EAGAIN.r   )	r4   r   r3   _Stream__scs_connectingr1   r2   r5   r6   r7   rn   r   r   r   r[   %  s   zStream.connectc              
   C   sD   z|  |W S  tjy! } ztj|dfW  Y d}~S d}~ww )a  Tries to receive up to 'n' bytes from this stream.  Returns a
        (error, string) tuple:

            - If successful, 'error' is zero and 'string' contains between 1
              and 'n' bytes of data.

            - On error, 'error' is a positive errno value.

            - If the connection has been closed in the normal fashion or if 'n'
              is 0, the tuple is (0, "").

        The recv function will not block waiting for data to arrive.  If no
        data have been received, it returns (errno.EAGAIN, "") immediately. N)_recvr   r7   r"   r>   get_exception_errnor8   nrV   r   r   r   recv6  s   zStream.recvc                 C   sR   |   }|dkr|dfS |dkrdS tjdkr!| jd u r!| |S d| j|fS )Nr   rt   r   rt   r   )r[   r   r   r   _Stream__recv_windowsry   )r8   rx   rq   r   r   r   ru   J  s   
zStream._recvc              
   C   sr  | j rLzt| j| jd}d| _ W n tjyK } z-|jtjjkr/d| _ t	j
dfW  Y d }~S |jtjv r<W Y d }~dS t	jdfW  Y d }~S d }~ww t| j|| j\}| _|rs|tjjkrhd| _ t	j
dfS |tjv rodS |dfS zt| j| jd}tj| jj W n% tjy } z|jtjv rW Y d }~dS |jdfW  Y d }~S d }~ww | jd | }dt|fS )NFTrt   rz   r   )_read_pendingr'   get_overlapped_resultr   r,   r*   r7   rL   ERROR_IO_INCOMPLETEr1   r2   pipe_disconnected_errorsEINVAL	read_file_read_bufferERROR_IO_PENDING
win32eventSetEventr.   bytes)r8   rx   
nBytesReadrV   errCode
recvBufferr   r   r   __recv_windowsV  sR   



zStream.__recv_windowsc              
   C   sB   z|  |W S  tjy  } ztj| W  Y d}~S d}~ww )aa  Tries to send 'buf' on this stream.

        If successful, returns the number of bytes sent, between 1 and
        len(buf).  0 is only a valid return value if len(buf) is 0.

        On error, returns a negative errno value.

        Will not block.  If no bytes can be immediately accepted for
        transmission, returns -errno.EAGAIN immediately.N)_sendr   r7   r"   r>   rv   r8   bufrV   r   r   r   send  s   zStream.sendc                 C   sd   |   }|dkr| S t|dkrdS t|tr|d}tjdkr,| jd u r,| |S | j	|S )Nr   zutf-8r   )
r[   rD   
isinstancestrencoder   r   r   _Stream__send_windowsr   )r8   r   rq   r   r   r   r     s   


zStream._sendc              
   C   s   | j rNzt| j| jd}d| _ W |S  tjyM } z.|jtjjkr/d| _	t
j W  Y d }~S |jtjv r?t
j W  Y d }~S t
j W  Y d }~S d }~ww t| j|| j\}}|rs|tjjkrhd| _ t
j S |ss|tjv rst
j S |S )NFT)_write_pendingr'   r}   r   r/   r*   r7   rL   r~   r|   r1   r2   r   
ECONNRESETr   
write_filer   )r8   r   nBytesWrittenrV   r   r   r   r   __send_windows  s:   
zStream.__send_windowsc                 C      d S r   r   rn   r   r   r   r^        z
Stream.runc                 C   r   r   r   r8   r_   r   r   r   ra     r   zStream.run_waitc                 C   s   |t jt jt jfv sJ | jt jkr|  d S | jt jkr!t j}tj	dkr.| 
|| d S |t jkr>|| jtjj d S || jtjj d S Nr   )r   	W_CONNECTW_RECVW_SENDr4   r6   immediate_waker3   r   r   _Stream__wait_windowsfd_waitr   r"   r_   POLLINPOLLOUT)r8   r_   waitr   r   r   r     s   

zStream.waitc              
   C   s  | j d urU|tjkrtjtjB tjB }tjj	}ntj
tjB tjB }tjj}zt| j | j| W n tjyK } ztd|j  W Y d }~nd }~ww || j| d S |tjkrk| jri|| jjtjj	 d S d S |tjkr| jr|| jjtjj d S d S |tjkrd S d S )Nz*failed to associate events with socket: %s)r   r   r   rO   FD_READ	FD_ACCEPTFD_CLOSEr"   r_   r   FD_WRITE
FD_CONNECTr   WSAEventSelectr0   r*   r7   rl   rW   strerrorr   r,   r.   r   r/   r   )r8   r_   r   maskeventrV   r   r   r   __wait_windows  sJ   





zStream.__wait_windowsc                 C      |  |tj d S r   )r   r   r   r   r   r   r   rb        zStream.connect_waitc                 C   r   r   )r   r   r   r   r   r   r   	recv_wait  r   zStream.recv_waitc                 C   r   r   )r   r   r   r   r   r   r   	send_wait  r   zStream.send_waitc                 C   sh   | j d ur
| j   | jd ur0| jrt| j | jjr#t| jj | jjr2t| jj d S d S d S r   )r   rG   r   r'   rk   r,   r.   r/   rn   r   r   r   __del__  s   


zStream.__del__c                 C   
   | t _d S r   )r   _SSL_private_key_file	file_namer   r   r   ssl_set_private_key_file     
zStream.ssl_set_private_key_filec                 C   r   r   )r   _SSL_certificate_filer   r   r   r   ssl_set_certificate_file  r   zStream.ssl_set_certificate_filec                 C   r   r   )r   _SSL_ca_cert_filer   r   r   r   ssl_set_ca_cert_file  r   zStream.ssl_set_ca_cert_fileNFr   )2__name__
__module____qualname____doc__r3   r5   r6   r   r   r   r   r   r   r   r/   r,   r   r|   ro   staticmethodr   r   r   r<   IPTOS_PREC_INTERNETCONTROLDSCP_DEFAULTr?   rF   rQ   rh   rG   rs   r[   ry   ru   r{   r   r   r   r^   ra   r   r   rb   r   r   r   r   r   r   r   r   r   r   r   6   sr    



!
B
*1

r   c                   @   sj   e Zd ZdZdZedd Zedd ZdddZed	d
 Z	dd Z
dd Zdd Zdd Zdd ZdS )r   NFc                 C   s   |  drdS dS )Npunix:FTr   r   r   r   r   r   &  r   zPassiveStream.needs_probesc                 C   s   |  d|  dB S )zReturns True if 'name' is a passive stream name in the form
        "TYPE:ARGS" and TYPE is a supported passive stream type (currently
        "punix:" or "ptcp"), otherwise False.r   ptcp:r   r   r   r   r   r   *  s   zPassiveStream.is_valid_namec                 C   sn   || _ || _|| _|d ur2t | _t | j_d| _	|
ddd }tjtjj|}t|| _|| _d S )NFr   r   )r	   r   r   r*   r+   r[   r'   r-   r.   connect_pendingr!   r"   r#   r$   r%   r&   r(   r)   	bind_path)r8   r@   r	   r   r   r;   r   r   r   r<   1  s   

zPassiveStream.__init__c              
   C   s  t | s
tjdfS | dd }| drptjtjj	|}t
jdkr6tjtjd|d\}}|r5|dfS ngz	t|d  W n
   tjdf Y S t|}t|dkrZtjdfS t|}|sftjdfS dt d| ||d	fS | d
rttjtj}|tjtjd | d}||d t|d f ntdz|d W n) tj y } zt!"d| t#$|j f  |  |j dfW  Y d}~S d}~ww dt || |fS )a  Attempts to start listening for remote stream connections.  'name'
        is a connection name in the form "TYPE:ARGS", where TYPE is an passive
        stream class's name and ARGS are stream class-specific. Currently the
        supported values for TYPE are "punix" and "ptcp".

        Returns (error, pstream): on success 'error' is 0 and 'pstream' is the
        new PassiveStream, on failure 'error' is a positive errno value and
        'pstream' is None.N   r   r   TwrA   r   r   r   r   r   r   zUnknown connection string
   z%s: listen: %s)%r   r   r1   rC   r   r"   r#   r$   r%   r&   r   r   r>   make_unix_socketr   SOCK_STREAMrF   rG   rE   r'   r(   rD   create_named_pipeAF_INET
setsockopt
SOL_SOCKETSO_REUSEADDRr!   bindint	Exceptionlistenr7   rl   rW   osr   )r	   r   r7   r@   rT   rU   remoterV   r   r   r   rF   ?  sL   











zPassiveStream.openc                 C   sf   | j dur
| j   | jdur t| jtj t| jjtj | j	dur1t
j| j	 d| _	dS dS )zCloses this PassiveStream.N)r   rG   r   r'   rk   rl   rm   r[   r.   r   r"   fatal_signalunlink_file_nowrn   r   r   r   rG   w  s   




zPassiveStream.closec              
   C   s   t jdkr| jdu r|  S 	 z6| j \}}tj| t jdkr3|jtj	kr3dt
|d| dfW S dt
|d|d t|d f dfW S  tjy~ } z,tj|}t jdkra|tjkratj}|tjkrptdt|  |dfW  Y d}~S d}~ww )	a  Tries to accept a new connection on this passive stream.  Returns
        (error, stream): if successful, 'error' is 0 and 'stream' is the new
        Stream object, and on failure 'error' is a positive errno value and
        'stream' is None.

        Will not block waiting for a connection.  If no connection is ready to
        be accepted, returns (errno.EAGAIN, None) immediately.r   NTr   zunix:%sz
ptcp:%s:%sr   z
accept: %s)r   r   r   _PassiveStream__accept_windowsacceptr"   r>   set_nonblockingfamilyAF_UNIXr   r   r7   rv   r1   r\   r2   rl   dbgr   r   )r8   r@   addrrV   r7   r   r   r   r     s,   

zPassiveStream.acceptc              
   C   sH  | j rHzt| j| jd W n6 tjyD } z)|jtjjkr,d| _ t	j
d fW  Y d }~S | jr5t| j t	jd fW  Y d }~S d }~ww d| _ t| j| j}|r~|tjjkr`d| _ t	j
d fS |tjjkrw| jrot| j d| _ t	jd fS t| jj t| j}|st	jd fS | j}|| _tj| jj dtd | jd|dfS )NFTr   r   )r   r'   r}   r   r[   r*   r7   rL   r~   r1   r2   rJ   rj   r   connect_named_piper   ERROR_PIPE_CONNECTEDr   r   r.   r   r)   rE   
ResetEventr   r	   )r8   rV   r7   rU   old_piper   r   r   __accept_windows  s>   	


zPassiveStream.__accept_windowsc                 C   sB   t jdks
| jd ur|| jtjj d S || jjtjj d S r   )	r   r   r   r   r"   r_   r   r[   r.   r   r   r   r   r     s   zPassiveStream.waitc                 C   sR   | j d ur
| j   | jd ur%| jrt| j | jjr't| jj d S d S d S r   )r   rG   r   r'   rk   _connectr.   r,   rn   r   r   r   r     s   


zPassiveStream.__del__r   )r   r   r   r[   r   r   r   r   r<   rF   rG   r   r   r   r   r   r   r   r   r   !  s    



7%r   c                 C   s   d| | f S )Na6  
Active %s connection methods:
  unix:FILE               Unix domain socket named FILE
  tcp:HOST:PORT           TCP socket to HOST with port no of PORT
  ssl:HOST:PORT           SSL socket to HOST with port no of PORT

Passive %s connection methods:
  punix:FILE              Listen on Unix domain socket FILEr   r   r   r   r   usage  s   r   c                   @   $   e Zd Zedd Zedd ZdS )
UnixStreamc                   C      dS r   r   r   r   r   r   r        zUnixStream.needs_probesc                 C   s   | }t jtjdd |S NT)r"   r>   r   r   r   )r;   rS   connect_pathr   r   r   rQ     s   
zUnixStream._openNr   r   r   r   r   rQ   r   r   r   r   r     
    
r   unixc                   @   r   )	TCPStreamc                   C   r   r   r   r   r   r   r   r     r   zTCPStream.needs_probesc              
   C   s   t jtj| d|\}}|s<z|tjtjd W ||fS  tjy; } z|	  t j
|d fW  Y d }~S d }~ww ||fS )Nr   r   )r"   r>   inet_open_activer   r   r   IPPROTO_TCPTCP_NODELAYr7   rG   rv   )r;   rS   r7   r@   rV   r   r   r   rQ     s   
zTCPStream._openNr   r   r   r   r   r     r   r   tcpc                       sd   e Zd Zedd Zedd Zedd Z fddZ fd	d
Z fddZ	 fddZ
  ZS )	SSLStreamc              
   C   s@   zt | W S  tjy } ztj|W  Y d }~S d }~ww r   )r   r?   sslSSLSyscallErrorr"   r>   rv   )r@   rV   r   r   r   r?     s   z%SSLStream.check_connection_completionc                   C   r   r   r   r   r   r   r   r     r   zSSLStream.needs_probesc           	   
   C   s
  t j| d}t jtj|\}}|d u r||fS ttj}tj	|_
| jtjO  _| jtjO  _|tj |tjtj |j|dd}t j||||}|sz|tjtjd W ||fS  tjy } z|  t j|d fW  Y d }~S d }~ww ||fS )Nr   F)do_handshake_on_connectr   )r"   r>   inet_parse_activeinet_create_socket_activer   r   r   
SSLContextPROTOCOL_SSLv23CERT_REQUIREDverify_modeoptionsOP_NO_SSLv2OP_NO_SSLv3load_verify_locationsr   r   load_cert_chainr   r   wrap_socketinet_connect_activer   r   r   r7   rG   rv   )	r;   rS   addressr   r@   ctxssl_sockr7   rV   r   r   r   rQ     s8   zSSLStream._openc                    s~   t t|  }|r|S z| j  W dS  tjy    tj Y S  tj	tj
tjtfy> } ztj|W  Y d }~S d }~ww )Nr   )superr   r[   r   do_handshaker   SSLWantReadErrorr1   r2   r   SSLZeroReturnErrorSSLEOFErrorOSErrorr"   r>   rv   rp   	__class__r   r   r[   0  s   

zSSLStream.connectc              
      s   z	t t| |W S  tjy   tjdf Y S  tjy2 } ztj	
|dfW  Y d }~S d }~w tjy<   Y dS  tjyV } ztj	
|dfW  Y d }~S d }~ww )Nrt   rz   )r  r   ru   r   r  r1   r2   r   r"   r>   rv   r  r   r7   rw   r  r   r   ry   A  s   zSSLStream.recvc              
      s   z	t t| |W S  tjy   tj  Y S  tjy0 } ztj	
| W  Y d }~S d }~w tjyI } ztj	
| W  Y d }~S d }~ww r   )r  r   r   r   SSLWantWriteErrorr1   r2   r   r"   r>   rv   r   r7   r   r  r   r   r   M  s   zSSLStream.sendc                    s<   | j rz	| j t j W n
 t jy   Y nw tt|  S r   )r   shutdown	SHUT_RDWRr7   r  r   rG   rn   r  r   r   rG   W  s   zSSLStream.close)r   r   r   r   r?   r   rQ   r[   ry   r   rG   __classcell__r   r   r  r   r     s    



r   r   )r1   r   r   r   
ovs.pollerr"   ovs.socket_utilovs.vlogr   ImportErrorr   ovs.winutilsr'   r*   r   rO   rJ   rl   Vlogr   objectr   r   r   r   r   r   r   r   r   r   r   <module>   sF   
   n 6[