î
&GäRIf  ã               @   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 m Z m Z m	 Z	 d  d l m
 Z
 m Z m Z d  d l m Z m Z d  d l m Z m Z m Z m Z d  d l m Z m Z m Z m Z m Z d  d l m Z m Z d  d l m Z m Z m Z d  d	 l m Z m  Z  d
 Z! d Z" d Z# d Z$ d Z% d Z& d Z' d Z( d  Z) Gd d „  d e ƒ Z* Gd d „  d e j+ ƒ Z, Gd d „  d ƒ Z- Gd d „  d e j+ ƒ Z. Gd d „  d e. ƒ Z/ Gd d „  d e. ƒ Z0 d S)!é    N)Ú
basestringÚbytesÚstr)ÚPropertyÚgetErrorMsgÚUID)ÚPackageQueueÚTinyPackageQueue)ÚPackageÚPACKAGE_HEARTBEATÚPACKAGE_CLOSEÚEINTR)Úcan_sendÚcan_recvÚsend_allÚrecv_allÚHEADER_SIZE)Ú
ConnectionÚTIMEOUT_MIN)ÚSTATUS_CLOSEDÚSTATUS_WAITINGÚSTATUS_HOSTING)ÚSTATUS_CONNECTEDÚSTATUS_CLOSINGzSocket error.z!Other end dropped the connection.zHandshake timed out.zHandshake failed.z2Handshake failed (context cannot connect to self).zClosed from other end.zLost track of the stream.zError in io thread.é
   i   c               @   s…   e  Z d  Z d Z d d d „ Z d d d „ Z d d	 d
 „ Z d d d „ Z d d „  Z d d d „ Z	 d d „  Z
 d d „  Z d S)ÚTcpConnectionar   TcpConnection(context, name='')
    
    The TcpConnection class implements a connection between two
    contexts that are in differenr processes or on different machines
    connected via the internet.
    
    This class handles the low-level communication for the context.    
    A ContextConnection instance wraps a single BSD socket for its 
    communication, and uses TCP/IP as the underlying communication 
    protocol. A persisten connection is used (the BSD sockets stay 
    connected). This allows to better distinguish between connection
    problems and timeouts caused by the other side being busy.
    
    Ú c             C   s>   d  |  _  d  |  _ t d | j Œ |  _ t j |  | | ƒ d  S)Né@   )Ú_sendingThreadÚ_receivingThreadr	   Z_queue_paramsÚ_qoutr   Ú__init__)ÚselfÚcontextÚname© r%   úO/Applications/pyzo2014a/lib/python3.4/site-packages/iep/yoton/connection_tcp.pyr!   :   s    		zTcpConnection.__init__Nc             C   s]  |  j  j ƒ  | d k	 r[ | j ƒ  \ |  _ |  _ | t k r[ | j ƒ  \ |  _ |  _ q[ n  t	 j
 |  | ƒ zÝ | t t g k rÑ | |  _ | j d ƒ t |  ƒ |  _ t |  ƒ |  _ |  j j ƒ  |  j j ƒ  n  | d k rGy |  j j ƒ  Wn t k
 rYn Xy |  j j ƒ  Wn t k
 r(Yn Xd |  _ d |  _ d |  _ n  Wd |  j  j ƒ  Xd S)a   _connected(status, bsd_socket=None)
        
        This method is called when a connection is made.
        
        Private method to apply the bsd_socket.
        Sets the socket and updates the status. 
        Also instantiates the IO threads.
        
        Né   r   )Ú_lockÚacquireÚgetsocknameÚ
_hostname1Ú_port1r   ÚgetpeernameÚ
_hostname2Ú_port2r   Ú_set_statusr   r   Ú_bsd_socketÚsetblockingÚSendingThreadr   ÚReceivingThreadr   ÚstartÚshutdownÚ	ExceptionÚcloseÚrelease)r"   ÚstatusÚ
bsd_socketr%   r%   r&   r0   G   s6    			zTcpConnection._set_statusr'   c             C   s>  t  j  t  j t  j ƒ } | j t  j t  j t ƒ | j t  j t  j t ƒ t j	 j
 d ƒ sx | j t  j t  j d ƒ n  xƒ t | | | ƒ D]H } y | j | | f ƒ PWqŒ t k
 rÓ | d k rÌ ‚  n  wŒ YqŒ XqŒ Wt | ƒ } d | d } t | ƒ ‚ | j d ƒ |  j t | ƒ t |  | ƒ |  _ |  j j ƒ  d S)z‹ Bind the bsd socket. Launches a dedicated thread that waits
        for incoming connections and to do the handshaking procedure.
        Úwinr'   zCould not bind to any of the z ports tried.r   N)ÚsocketÚAF_INETÚSOCK_STREAMÚ
setsockoptÚ
SOL_SOCKETÚ	SO_SNDBUFÚSOCKET_BUFFERS_SIZEÚ	SO_RCVBUFÚsysÚplatformÚ
startswithÚSO_REUSEADDRÚrangeÚbindr7   r   ÚIOErrorÚlistenr0   r   Ú
HostThreadZ_hostThreadr5   )r"   ÚhostnameÚportÚ	max_triesÚsÚport2Útmpr%   r%   r&   Ú_bind‹   s(    zTcpConnection._bindg      ð?c             C   s´  t  j  t  j t  j ƒ } | j t  j t  j t ƒ | j t  j t  j t ƒ | d k r_ d } n  d } t j ƒ  | } xw | rî t j ƒ  | k  rî y | j	 | | f ƒ d } Wn) t  j
 k
 rÅ Yn t  j k
 rÙ Yn Xt j | d ƒ qx W| s5t j ƒ  \ } } }	 ~	 t | ƒ }
 t d | | |
 f ƒ ‚ n  t | ƒ } | j |  j ƒ \ } } | sŽ|  j d ƒ | s{d } n  t d | ƒ ‚ n  | \ |  _ |  _ |  j t | ƒ d	 S)
z$ Connect to a bound socket.
        g{®Gáz„?FTg      Y@zCannot connect to %s on %i: %sr   zproblem during handshakezCould not connect: N)r=   r>   r?   r@   rA   rB   rC   rD   ÚtimeÚconnectÚerrorÚtimeoutÚsleeprE   Úexc_infor   rK   Ú
HandShakerÚshake_hands_as_clientÚid1r0   Ú_id2Ú_pid2r   )r"   rN   rO   rX   rQ   ÚokÚ	timestampÚtypeÚvalueÚtbÚerrÚhÚsuccessÚinfor%   r%   r&   Ú_connectÄ   s<    	
	zTcpConnection._connectc             C   s:   |  j  j t ƒ y |  j d ƒ Wn t k
 r5 Yn Xd S)z Send PACKAGE_CLOSE.
        g      ð?N)r    Úpushr   Úflushr7   )r"   r%   r%   r&   Ú_notify_other_end_of_closingù   s
    z*TcpConnection._notify_other_end_of_closingg      @c             C   sˆ   |  j  s t d ƒ n  |  j j t ƒ t j ƒ  | } xK |  j  rƒ |  j j ƒ  rƒ t j d ƒ t j ƒ  | k r9 t d ƒ ‚ q9 q9 Wd S)zn Put a dummy message on the queue and spinlock until
        the thread has popped it from the queue.
        zCannot flush if not connected.g{®Gáz„?zSending the packages timed out.N)Úis_connectedÚRuntimeErrorr    rj   r   rU   ÚemptyrY   )r"   rX   ra   r%   r%   r&   Ú_flush  s    	zTcpConnection._flushc             C   s   |  j  j | ƒ d S)zU Put package on the queue, where the sending thread will
        pick it up.
        N)r    rj   )r"   Úpackager%   r%   r&   Ú_send_package  s    zTcpConnection._send_packagec             C   s   |  j  j j | ƒ d S)z> Put package in queue, but bypass potential blocking.
        N)r    Ú_qÚappend)r"   rq   r%   r%   r&   Ú_inject_package  s    zTcpConnection._inject_package)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r!   r0   rT   ri   rl   rp   rr   ru   r%   r%   r%   r&   r   *   s   D95
r   c               @   s:   e  Z d  Z d Z d d „  Z d d „  Z d d „  Z d S)	rM   a€   HostThread(context_connection, bds_socket)
    
    The host thread is used by the ContextConnection when hosting a 
    connection. This thread waits for another context to connect 
    to it, and then performs the handshaking procedure.
    
    When a successful connection is made, the context_connection's 
    _connected() method is called and this thread then exits.
    
    c             C   s3   t  j j |  ƒ | |  _ | |  _ |  j d ƒ d  S)NT)Ú	threadingÚThreadr!   Ú_context_connectionÚ_bsd_host_socketÚ	setDaemon)r"   Úcontext_connectionr;   r%   r%   r&   r!   2  s    		zHostThread.__init__c             C   sÍ   xº |  j  j r¼ |  j ƒ  } | s' q n  |  j  j s7 Pn  t | ƒ } | j |  j  j ƒ \ } } | r‡ | d |  j  _ | d |  j  _ n t d | ƒ q |  j	 j
 ƒ  |  j  j t | ƒ Pq W|  `  |  `	 d S)z€ run()
        
        The main loop. Waits for a connection and performs handshaking
        if successfull.
        
        r   r'   zYoton: Handshake failed: N)r|   Ú
is_waitingÚ_wait_for_connectionr[   Úshake_hands_as_hostr]   r^   r_   Úprintr}   r8   r0   r   )r"   rQ   Zhsrg   rh   r%   r%   r&   Úrun=  s$    	zHostThread.runc             C   s�   |  j  j d ƒ x† |  j j r˜ y |  j  j ƒ  \ } } | SWq t j k
 rS Yq t j k
 r” t j	 ƒ  \ } } } ~ | j
 t k r� ‚  n  Yq Xq Wd S)zµ _wait_for_connection()    
            
        The thread will wait here until someone connects. When a 
        connections is made, the new socket is returned.
        
        g      Ð?N)r}   Ú
settimeoutr|   r€   Úacceptr=   rX   rW   rE   rZ   Úerrnor   )r"   rQ   Úaddrrb   rc   rd   r%   r%   r&   r�   h  s    	zHostThread._wait_for_connectionN)rv   rw   rx   ry   r!   r„   r�   r%   r%   r%   r&   rM   &  s   
+rM   c               @   sU   e  Z d  Z d Z d d „  Z d d „  Z d d „  Z d d	 d
 „ Z d d „  Z d S)r[   a‹   HandShaker(bsd_socket)
    
    Class that performs the handshaking procedure for Tcp connections.
    
    Essentially, the connecting side starts by sending 'YOTON!' 
    followed by its id as a hex string. The hosting side responds
    with the same message (but with a different id).
    
    This process is very similar to a client/server pattern (both 
    messages are also terminated with '
'). This is done such that
    if for example a web client tries to connect, a sensible error
    message can be returned. Or when a ContextConnection tries to connect
    to a web server, it will be able to determine the error gracefully.
    
    c             C   s   | |  _  d  S)N)r1   )r"   r;   r%   r%   r&   r!   ”  s    zHandShaker.__init__c       
      C   s   d t  | ƒ j ƒ  t j ƒ  f } |  j ƒ  } | s> d t f S| j d ƒ ryT | d d … j d d ƒ } | d | d } } t | d	 ƒ t | d
 ƒ } } Wn) t	 k
 rÌ |  j
 d ƒ d t f SYn X|  j
 | ƒ }	 | | k rò d t f Sd | | f f Sn |  j
 d ƒ d t f Sd S)aU   _shake_hands_as_host(id)
        
        As the host, we wait for the client to ask stuff, so when
        for example a http client connects, we can stop the connection.
        
        Returns (success, info), where info is the id of the context at
        the other end, or the error message in case success is False.
        
        zYOTON!%s.%iFzYOTON!é   NÚ.r'   r   é   r   zERROR: could not parse id.TzERROR: this is Yoton.)r   Úget_hexÚosÚgetpidÚ_recv_during_handshakingÚSTOP_HANDSHAKE_TIMEOUTrG   ÚsplitÚintr7   Ú_send_during_handshakingÚSTOP_HANDSHAKE_FAILEDÚSTOP_HANDSHAKE_SELF)
r"   ÚidÚmessageZrequestrS   Úid2_strÚpid2_strÚid2Úpid2rW   r%   r%   r&   r‚   š  s$    "
#
zHandShaker.shake_hands_as_hostc       
      C   s  d t  | ƒ j ƒ  t j ƒ  f } |  j | ƒ } |  j ƒ  } | sM d t f S| j d ƒ rø yT | d d … j d d ƒ } | d | d } } t	 | d	 ƒ t	 | d
 ƒ } }	 Wn t
 k
 rÎ d t f SYn X| | k rå d t f Sd | |	 f f Sn
 d t f Sd S)aU   _shake_hands_as_client(id)
        
        As the client, we ask the host whether it is a Yoton context
        and whether the channels we want to support are all right.
        
        Returns (success, info), where info is the id of the context at
        the other end, or the error message in case success is False.
        
        zYOTON!%s.%iFzYOTON!r‰   NrŠ   r'   r   r‹   r   T)r   rŒ   r�   rŽ   r“   r�   r�   rG   r‘   r’   r7   r”   r•   )
r"   r–   r—   rW   ZresponserS   r˜   r™   rš   r›   r%   r%   r&   r\   Â  s     "
#
z HandShaker.shake_hands_as_clientFc             C   s   t  |  j | d | ƒ S)Nz
)r   r1   )r"   Útextr6   r%   r%   r&   r“   é  s    z#HandShaker._send_during_handshakingc             C   s   t  |  j d d ƒ S)Ng       @T)r   r1   )r"   r%   r%   r&   r�   í  s    z#HandShaker._recv_during_handshakingN)	rv   rw   rx   ry   r!   r‚   r\   r“   r�   r%   r%   r%   r&   r[   ƒ  s   ('r[   c               @   s:   e  Z d  Z d Z d d „  Z d d „  Z d d „  Z d S)	ÚBaseIOThreadzh The base class for the sending and receiving IO threads.
    Implements some common functionality.
    c             C   s6   t  j j |  ƒ |  j d ƒ | |  _ | j |  _ d  S)NT)rz   r{   r!   r~   r|   r1   )r"   r   r%   r%   r&   r!   ÷  s    	zBaseIOThread.__init__c             C   sK   |  j  } |  j } |  `  |  ` y |  j | | ƒ Wn t k
 rF Yn Xd S)z† Method to prepare to enter main loop. There is a try-except here
        to catch exceptions caused by interpreter shutdown.
        N)r|   r1   Úrun2r7   )r"   r   r;   r%   r%   r&   r„     s    		zBaseIOThread.runc             C   sÛ   d |  j  j } d d „  } y, |  j | | ƒ } | rG | j | ƒ n  WnŒ t j k
 rƒ t t ƒ  } | j d | | f ƒ YnT t k
 rÖ t ƒ  } t	 | } | j d | | f ƒ | d | ƒ | | ƒ Yn Xd S)z° Method to enter main loop. There is a try-except here to
        catch exceptions in the main loop (such as socket errors and 
        errors due to bugs in the code.
        zyoton.c             S   s+   t  j j t |  ƒ d ƒ t  j j ƒ  d  S)NÚ
)rE   Ú
__stderr__Úwriter   rk   )re   r%   r%   r&   ÚwriteErr#  s    z#BaseIOThread.run2.<locals>.writeErrz%s, %szException in %s.N)
Ú	__class__rv   Ú_runÚclose_on_problemr=   rW   ÚSTOP_SOCKET_ERRORr   r7   ÚSTOP_THREAD_ERROR)r"   r   r;   Z	classNamer¢   Zstop_reasonÚmsgZerrmsgr%   r%   r&   rž     s     	
zBaseIOThread.run2N)rv   rw   rx   ry   r!   r„   rž   r%   r%   r%   r&   r�   ò  s   r�   c               @   s"   e  Z d  Z d Z d d „  Z d S)r3   zÇ The thread that reads packages from the queue and sends them over
    the socket. It uses a timeout while reading from the queue, so it can
    send heart beat packages if no packages are send.
    c       	      C   s›   d t  } | j } | j } | j } xo t j d ƒ y | j | ƒ } Wn( | j k
 rr t } | j	 sn d SYn Xx | j
 ƒ  D] } | | ƒ q€ Wq( d S)zI  The main loop. Get package from queue, send package to socket.
        g      à?r   N)r   r    ÚsendÚsendallrU   rY   ÚpopÚEmptyr   rm   Úparts)	r"   r   r;   rX   ZqueueZsocket_sendZsocket_sendallrq   Úpartr%   r%   r&   r¤   H  s    
					zSendingThread._runN)rv   rw   rx   ry   r¤   r%   r%   r%   r&   r3   B  s   r3   c               @   s:   e  Z d  Z d Z d d „  Z d d „  Z d d „  Z d S)	r4   a   The thread that reads packages from the socket and passes them to
    the kernel. It uses select() to see if data is available on the socket.
    This allows using a timeout without putting the socket in timeout mode.
    
    If the timeout has expired, the timedout event for the connection is
    emitted.
    
    Upon receiving a package, the _recv_package() method of the context
    is called, so this thread will eventually dispose the package in 
    one or more queues (of the channel or of another connection).
    
    c       
      C   s]  | j  } | j j } t j } t } d } x,t j d ƒ x¢ y t | | j	 ƒ } Wn( t
 k
 r} t j d t ƒ  ƒ ‚ Yn X| rª | r¦ d } | j j | d ƒ n  Pq= | j s· d S| s= d } | j j | d ƒ q= q= q= |  j | | | ƒ }	 |	 d k rq- n t |	 t ƒ r|	 Sy | |	 | ƒ Wq- t
 k
 rUt d ƒ t t ƒ  ƒ Yq- Xq- d S)zN The main loop. Get package from socket, deposit package in queue(s).
        Fr   z
select(): NTz,Error depositing package in ReceivingThread.)ÚrecvÚ_contextZ_recv_packager
   Úfrom_headerr   rU   rY   r   Ú_timeoutr7   r=   rW   r   ÚtimedoutÚemitrm   Ú_getPackageÚ
isinstancer   rƒ   )
r"   r   r;   Úsocket_recvZrecv_packageÚpackage_from_headerÚHSZtimedOutr`   rq   r%   r%   r&   r¤   s  sB    			
zReceivingThread._runc             C   sº   y |  j  | | ƒ } Wn t k
 r. t SYn X| | ƒ \ } } | sK t S| d k r€ | j d k ri n | j d k r| t Sd Sy |  j  | | ƒ | _ Wn t k
 r± t SYn X| Sd S)z< Get exactly one package from the socket. Blocking.
        r   r'   N)Ú_recv_n_bytesÚEOFErrorÚSTOP_EOFÚSTOP_LOST_TRACKÚ_source_seqÚSTOP_CLOSED_FROM_THEREÚ_data)r"   r·   r¹   r¸   Úheaderrq   Úsizer%   r%   r&   rµ   ¬  s$    		zReceivingThread._getPackagec             C   s™   | | ƒ } t  | ƒ d k r* t ƒ  ‚ n  | t  | ƒ 8} | d k rJ | S| g } x3 | rˆ | | ƒ } | j | ƒ | t  | ƒ 8} qV Wt ƒ  j | ƒ S)z2 Receive exactly n bytes from the socket.
        r   )Úlenr»   rt   r   Újoin)r"   r·   ÚnÚdatar­   r%   r%   r&   rº   Í  s    		zReceivingThread._recv_n_bytesN)rv   rw   rx   ry   r¤   rµ   rº   r%   r%   r%   r&   r4   e  s   9!r4   i (  )1r�   rE   rU   r=   rz   ÚyotonÚ
yoton.miscr   r   r   r   r   r   r   r	   Ú
yoton.corer
   r   r   r   r   r   r   r   r   Úyoton.connectionr   r   r   r   r   r   r   r¦   r¼   r�   r”   r•   r¿   r½   r§   rC   r   r{   rM   r[   r�   r3   r4   r%   r%   r%   r&   Ú<module>   s8   "(ü]oP#