î
&GäRP  ã               @   s  d  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	 m
 Z
 m Z m Z m Z d d l m Z m Z m Z m Z d d l m Z m Z d d l m Z d d	 l m Z m Z d d
 l m Z d d l m Z Gd d „  d e ƒ Z d S)z3 Module yoton.context

Defines the context class.

é    N)Úcore)Ú
connection)Ú
basestringÚbytesÚstrÚxrangeÚsplit_address)ÚPropertyÚgetErrorMsgÚUIDÚPackageQueue)ÚPackageÚBUF_MAX_LEN)ÚSLOT_CONTEXT)ÚConnectionCollectionÚ
Connection)ÚTcpConnection)ÚItcConnectionc               @   s	  e  Z d  Z d Z d d d d „ Z d d „  Z d d	 „  Z e d
 d „  ƒ Z e d d „  ƒ Z	 e d d „  ƒ Z
 e d d „  ƒ Z d d d d „ Z d d d d „ Z d d d „ Z d d d „ Z d d d „ Z d  d! „  Z d" d# „  Z d$ d% „  Z d& d' „  Z d S)(ÚContexta­   Context(verbose=0, queue_params=None)
    
    A context represents a node in the network. It can connect to 
    multiple other contexts (using a yoton.Connection. 
    These other contexts can be in 
    another process on the same machine, or on another machine
    connected via a network or the internet.
    
    This class represents a context that can be used by channel instances
    to communicate to other channels in the network. (Thus the name.)
    
    The context is the entity that queue routes the packages produced 
    by the channels to the other context in the network, where
    the packages are distributed to the right channels. A context queues
    packages while it is not connected to any other context.
    
    If messages are send on a channel registered at this context while
    the context is not connected, the messages are stored by the
    context and will be send to the first connecting context.
    
    Example 1
    ---------
    # Create context and bind to a port on localhost
    context = yoton.Context()
    context.bind('localhost:11111')
    # Create a channel and send a message
    pub = yoton.PubChannel(context, 'test')
    pub.send('Hello world!')
    
    Example 2
    ---------
    # Create context and connect to the port on localhost
    context = yoton.Context()
    context.connect('localhost:11111')
    # Create a channel and receive a message
    sub = yoton.SubChannel(context, 'test')
    print(sub.recv() # Will print 'Hello world!'
    
    Queue params
    ------------
    The queue_params parameter allows one to specify the package queues
    used in the system. It is recommended to use the same parameters
    for every context in the network. The value of queue_params should
    be a 2-element tuple specifying queue size and discard mode. The
    latter can be 'old' (default) or 'new', meaning that if the queue
    is full, either the oldest or newest messages are discarted.
    
    r   Nc             C   sÇ   | |  _  | d  k r$ t d f } n  t | t ƒ oB t | ƒ d k sT t d ƒ ‚ n  | |  _ t ƒ  j ƒ  |  _	 i  |  _
 i  |  _ g  |  _ t j ƒ  |  _ t | Œ  |  _ d |  _ d |  _ i  |  _ d  S)NÚoldé   z)queue_params should be a 2-element tuple.r   )Ú_verboser   Ú
isinstanceÚtupleÚlenÚ
ValueErrorÚ_queue_paramsr   Úget_intÚ_idÚ_sending_channelsÚ_receiving_channelsÚ_connectionsÚ	threadingÚRLockÚ_connections_lockr   Ú_startupQueueÚ	_send_seqÚ	_recv_seqÚ_source_map)ÚselfÚverboseZqueue_params© r+   úH/Applications/pyzo2014a/lib/python3.4/site-packages/iep/yoton/context.pyÚ__init__O   s    	!						zContext.__init__c             C   s/   x |  j  D] } | j d ƒ q
 W|  j ƒ  d S)a©   close()
        
        Close the context in a nice way, by closing all connections
        and all channels.
        
        Closing a connection means disconnecting two contexts. Closing
        a channel means disasociating a channel from its context. 
        Unlike connections and channels, a Context instance can be reused 
        after closing (although this might not always the best strategy).
        
        zClosed by the context.N)Úconnections_allÚcloseÚclose_channels)r)   Úcr+   r+   r,   r/   p   s    zContext.closec             C   sa   d d „  |  j  j ƒ  Dƒ } d d „  |  j j ƒ  Dƒ } x" t | | ƒ D] } | j ƒ  qI Wd S)z¤ close_channels()
        
        Close all channels associated with this context. This does
        not close the connections. See also close().
        
        c             S   s   g  |  ] } | ‘ q Sr+   r+   )Ú.0r1   r+   r+   r,   ú
<listcomp>Ž   s   	 z*Context.close_channels.<locals>.<listcomp>c             S   s   g  |  ] } | ‘ q Sr+   r+   )r2   r1   r+   r+   r,   r3   �   s   	 N)r   Úvaluesr    Úsetr/   )r)   Z	channels1Z	channels2r1   r+   r+   r,   r0   …   s    	zContext.close_channelsc          
   C   s:   |  j  j ƒ  z d d „  |  j Dƒ SWd |  j  j ƒ  Xd S)a6   Get a list of all Connection instances currently
        associated with this context, including pending connections 
        (connections waiting for another end to connect).
        In addition to normal list indexing, the connections objects can be
        queried from this list using their name.
        c             S   s   g  |  ] } | j  r | ‘ q Sr+   )Úis_alive)r2   r1   r+   r+   r,   r3   ¤   s   	 z+Context.connections_all.<locals>.<listcomp>N)r$   Úacquirer!   Úrelease)r)   r+   r+   r,   r.   ˜   s    
zContext.connections_allc          
   C   s    |  j  j ƒ  z~ t ƒ  } g  } xC |  j D]8 } | j sH | j | ƒ q) | j r) | j | ƒ q) q) Wx | D] } |  j j | ƒ ql W| SWd |  j  j ƒ  Xd S)zÚ Get a list of the Connection instances currently
        active for this context. 
        In addition to normal list indexing, the connections objects can be
        queried  from this list using their name.
        N)	r$   r7   r   r!   r6   ÚappendÚis_connectedÚremover8   )r)   ÚcopyZ	to_remover1   r+   r+   r,   Úconnections©   s    			zContext.connectionsc             C   s   t  |  j ƒ S)z‹ Get the number of connected contexts. Can be used as a boolean
        to check if the context is connected to any other context.
        )r   r=   )r)   r+   r+   r,   Úconnection_countÉ   s    zContext.connection_countc             C   s   |  j  S)z) The 8-byte UID of this context.
        )r   )r)   r+   r+   r,   ÚidÑ   s    z
Context.idé   Ú c          
   C   s    |  j  t | ƒ \ } } } t |  | ƒ } | j | | | ƒ |  j j ƒ  z@ x) t |  j ƒ ry | j |  j j	 ƒ  ƒ qQ W|  j
 j | ƒ Wd |  j j ƒ  X| S)a¸   bind(address, max_tries=1, name='')
        
        Setup a connection with another Context, by being the host.
        This method starts a thread that waits for incoming connections.
        Error messages are printed when an attemped connect fails. the
        thread keeps trying until a successful connection is made, or until
        the connection is closed.
        
        Returns a Connection instance that represents the
        connection to the other context. These connection objects 
        can also be obtained via the Context.connections property.
        
        Parameters
        ----------
        address : str
            Should be of the shape hostname:port. The port should be an
            integer number between 1024 and 2**16. If port does not 
            represent a number, a valid port number is created using a 
            hash function.
        max_tries : int
            The number of ports to try; starting from the given port, 
            subsequent ports are tried until a free port is available. 
            The final port can be obtained using the 'port' property of
            the returned Connection instance.
        name : string
            The name for the created Connection instance. It can
            be used as a key in the connections property.
        
        Notes on hostname
        -----------------
        The hostname can be:
          * The IP address, or the string hostname of this computer. 
          * 'localhost': the connections is only visible from this computer. 
            Also some low level networking layers are bypassed, which results
            in a faster connection. The other context should also connect to
            'localhost'.
          * 'publichost': the connection is visible by other computers on the 
            same network. Optionally an integer index can be appended if
            the machine has multiple IP addresses (see socket.gethostbyname_ex).
        
        N)r=   r   r   Ú_bindr$   r7   r   r%   Ú_inject_packageÚpopr!   r9   r8   )r)   ÚaddressÚ	max_triesÚnameÚprotocolÚhostnameÚportr   r+   r+   r,   ÚbindÛ   s    ,zContext.bindg      ð?c       
      C   s×   |  j  t | ƒ \ } } } t |  | ƒ } | j | | | ƒ |  j j ƒ  z: x# |  j rs | j |  j j ƒ  ƒ qQ W|  j	 j
 | ƒ Wd |  j j ƒ  Xd j d ƒ } t | t |  j d d d d ƒ }	 |  j |	 ƒ | S)a   connect(self, address, timeout=1.0, name='')
        
        Setup a connection with another context, by connection to a 
        hosting context. An error is raised when the connection could
        not be made.
        
        Returns a Connection instance that represents the
        connection to the other context. These connection objects 
        can also be obtained via the Context.connections property.
        
        Parameters
        ----------
        address : str
            Should be of the shape hostname:port. The port should be an
            integer number between 1024 and 2**16. If port does not 
            represent a number, a valid port number is created using a 
            hash function.
        max_tries : int
            The number of ports to try; starting from the given port, 
            subsequent ports are tried until a free port is available. 
            The final port can be obtained using the 'port' property of
            the returned Connection instance.
        name : string
            The name for the created Connection instance. It can
            be used as a key in the connections property.
        
        Notes on hostname
        -----------------
        The hostname can be:
          * The IP address, or the string hostname of this computer. 
          * 'localhost': the connection is only visible from this computer. 
            Also some low level networking layers are bypassed, which results
            in a faster connection. The other context should also host as
            'localhost'.
          * 'publichost': the connection is visible by other computers on the 
            same network. Optionally an integer index can be appended if
            the machine has multiple IP addresses (see socket.gethostbyname_ex).
        
        NÚNEW_CONNECTIONzutf-8r   )r=   r   r   Ú_connectr$   r7   r%   rC   rD   r!   r9   r8   Úencoder   r   r   Ú_send_package)
r)   rE   ÚtimeoutrG   rH   rI   rJ   r   ÚbbÚpr+   r+   r,   Úconnect$  s    *!zContext.connectg      @c             C   s%   x |  j  D] } | j | ƒ q
 Wd S)a>   flush(timeout=5.0)
        
        Wait until all pending messages are send. This will flush all
        messages posted from the calling thread. However, it is not
        guaranteed that no new messages are posted from another thread.
        
        Raises an error when the flushing times out.
        
        T)r=   Úflush)r)   rP   r1   r+   r+   r,   rT   p  s    zContext.flushc             C   s9   | |  j  k r( t d t | ƒ ƒ ‚ n  | |  j  | <d S)z³ _register_sending_channel(channel, slot, slotname='')
        
        The channel objects use this method to register themselves 
        at a particular slot.
        
        zSlot not free: N)r   r   r   )r)   ÚchannelÚslotÚslotnamer+   r+   r,   Ú_register_sending_channel…  s    	z!Context._register_sending_channelc             C   s9   | |  j  k r( t d t | ƒ ƒ ‚ n  | |  j  | <d S)zµ _register_receiving_channel(channel, slot, slotname='')
        
        The channel objects use this method to register themselves 
        at a particular slot.
        
        zSlot not free: N)r    r   r   )r)   rU   rV   rW   r+   r+   r,   Ú_register_receiving_channel•  s    	z#Context._register_receiving_channelc             C   se   x^ |  j  |  j g D]J } xA d d „  | j ƒ  Dƒ D]& } | | | k r3 | j | ƒ q3 q3 Wq Wd S)z¸ _unregister_channel(channel)
        
        Unregisters the given channel. That channel can no longer
        receive messages, and should no longer send messages.
        
        c             S   s   g  |  ] } | ‘ q Sr+   r+   )r2   Úkeyr+   r+   r,   r3   ­  s   	 z/Context._unregister_channel.<locals>.<listcomp>N)r    r   ÚkeysrD   )r)   rU   ÚDrZ   r+   r+   r,   Ú_unregister_channel¥  s     zContext._unregister_channelc          
   C   s“   |  j  d 7_  |  j  | _ |  j j ƒ  zV d } x0 |  j D]% } | j r; | j | ƒ d } q; q; W| s} |  j j | ƒ n  Wd |  j j	 ƒ  Xd S)a   _send_package(package)
        
        Used by the channels to send a package into the network.
        This method routes the package to all currentlt connected
        connections. If there are none, the packages is queued at
        the context.
        
        r@   FTN)
r&   Ú_source_seqr$   r7   r!   r6   rO   r%   Úpushr8   )r)   ÚpackageÚokr1   r+   r+   r,   rO   µ  s    	zContext._send_packagec       	      C   s‰  | j  } d } d } |  j j | j d ƒ } | | j k  rÇ | j |  j | j <| j d k rˆ |  j d 7_ |  j | _ d \ } } qÇ | j |  j k r¾ |  j d 7_ |  j | _ d } qÇ d } n  | r/|  j j	 ƒ  zA x: |  j
 D]/ } | | k sç | j r	qç n  | j | ƒ qç WWd |  j j ƒ  Xn  | r…| t k rQ|  j | ƒ q…|  j j | d ƒ } | d k	 r…| j | ƒ q…n  d S)a(   _recv_package(package, connection)
        
        Used by the connections to receive a package at this
        context. The package is distributed to all connections
        except the calling one. The package is also distributed
        to the right channel (if applicable).
        
        Fr   r@   TN)TT)Ú_slotr(   ÚgetÚ
_source_idr^   Ú_dest_idr'   r   r$   r7   r!   r6   rO   r8   r   Ú_recv_context_packager    Ú_recv_package)	r)   r`   r   rV   Zsend_furtherZdeposit_hereZlast_seqr1   rU   r+   r+   r,   rg   Ò  s:    			zContext._recv_packagec             C   sî   | j  j d ƒ } | d k rˆ |  j j ƒ  zI xB |  j D]7 } | j r8 | j | j k r8 | j t	 j
 d ƒ q8 q8 WWd |  j j ƒ  Xnb | d k rÜ xS |  j j ƒ  D]1 } t | d ƒ r¤ t | d ƒ r¤ | j ƒ  q¤ q¤ Wn t d | ƒ d S)	z¼ _recv_context_package(package)
        
        Process a package addressed at the context itself. This is how
        the context handles higher-level connection tasks.
        
        zutf-8ZCLOSE_CONNECTIONFNrL   Z_current_messageÚ	send_lastz)Yoton: Received unknown context message: )Ú_dataÚdecoder$   r7   r=   r:   Úid2rd   r/   r   ÚSTOP_CLOSED_FROM_THEREr8   r   r4   Úhasattrrh   Úprint)r)   r`   Úmessager1   rU   r+   r+   r,   rf     s    	zContext._recv_context_package)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r-   r/   r0   Úpropertyr.   r=   r>   r?   rK   rS   rT   rX   rY   r]   rO   rg   rf   r+   r+   r+   r,   r      s"   0! 
IL>r   )rs   ÚosÚsysÚtimer"   Úyotonr   r   Ú
yoton.miscr   r   r   r   r   r	   r
   r   r   Ú
yoton.corer   r   r   Úyoton.connectionr   r   Úyoton.connection_tcpr   Zyoton.connection_itcr   Úobjectr   r+   r+   r+   r,   Ú<module>   s   $("