3
ñ\ð]b  ã               @   sb   d dl Z d dlZyd dljjZW n ek
r8   dZY nX d dlZddlmZ G dd„ deƒZ	dS )é    Né   )ÚPubSubManagerc                   s>   e Zd ZdZdZd‡ fdd„	Zd	d
„ Zdd„ Zdd„ Z‡  Z	S )Ú
ZmqManagera?  zmq based client manager.

    NOTE: this zmq implementation should be considered experimental at this
    time. At this time, eventlet is required to use zmq.

    This class implements a zmq backend for event sharing across multiple
    processes. To use a zmq backend, initialize the :class:`Server` instance as
    follows::

        url = 'zmq+tcp://hostname:port1+port2'
        server = socketio.Server(client_manager=socketio.ZmqManager(url))

    :param url: The connection URL for the zmq message broker,
                which will need to be provided and running.
    :param channel: The channel name on which the server sends and receives
                    notifications. Must be the same in all the servers.
    :param write_only: If set to ``True``, only initialize to emit events. The
                       default of ``False`` initializes the class for emitting
                       and receiving.

    A zmq message broker must be running for the zmq_manager to work.
    you can write your own or adapt one from the following simple broker
    below::

        import zmq

        receiver = zmq.Context().socket(zmq.PULL)
        receiver.bind("tcp://*:5555")

        publisher = zmq.Context().socket(zmq.PUB)
        publisher.bind("tcp://*:5556")

        while True:
            publisher.send(receiver.recv())
    Úzmqúzmq+tcp://localhost:5555+5556ÚsocketioFNc                sÜ   t d krtdƒ‚tjdƒ}|jdƒo,|j|ƒs:td| ƒ‚|jddƒ}|jdƒ\}}|jdƒd }|j||ƒ}	t jƒ j	t j
ƒ}
|
j|ƒ t jƒ j	t jƒ}|jt jdƒ |j|	ƒ |
| _|| _|| _tt| ƒj|||d
� d S )NzJzmq package is not installed (Run "pip install pyzmq" in your virtualenv).z
:\d+\+\d+$z
zmq+tcp://zunexpected connection string: zzmq+Ú ú+ú:r   )ÚchannelÚ
write_onlyÚloggeréÿÿÿÿ)r   ÚRuntimeErrorÚreÚcompileÚ
startswithÚsearchÚreplaceÚsplitÚContextÚsocketZPUSHÚconnectZSUBZsetsockopt_stringZ	SUBSCRIBEÚsinkÚsubr   Úsuperr   Ú__init__)ÚselfÚurlr   r   r   ÚrZsink_urlZsub_portZ	sink_portZsub_urlr   r   )Ú	__class__© ú?/tmp/pip-build-mqc4i71p/python-socketio/socketio/zmq_manager.pyr   3   s(    


zZmqManager.__init__c             C   s    t jd| j|dœƒ}| jj|ƒS )NÚmessage)Útyper   Údata)ÚpickleÚdumpsr   r   Úsend)r   r%   Zpickled_datar!   r!   r"   Ú_publishS   s
    
zZmqManager._publishc             c   s"   x| j jƒ }|d k	r|V  qW d S )N)r   Úrecv)r   Úresponser!   r!   r"   Ú
zmq_listen]   s    
zZmqManager.zmq_listenc             c   s|   xv| j ƒ D ]j}t|tjƒr>ytj|ƒ}W n tk
r<   Y nX t|tƒr
|d dkr
|d | jkr
d|kr
|d V  q
W d S )Nr$   r#   r   r%   )	r,   Ú
isinstanceÚsixÚbinary_typer&   ÚloadsÚ	ExceptionÚdictr   )r   r#   r!   r!   r"   Ú_listenc   s    
zZmqManager._listen)r   r   FN)
Ú__name__Ú
__module__Ú__qualname__Ú__doc__Únamer   r)   r,   r3   Ú__classcell__r!   r!   )r    r"   r      s   #   
r   )
r&   r   Zeventlet.green.zmqZgreenr   ÚImportErrorr.   Zpubsub_managerr   r   r!   r!   r!   r"   Ú<module>   s   
