3
ñ\ð]^  ã               @   sh   d dl Z d dlZd dlZyd dlZW n ek
r<   dZY nX ddlmZ e jdƒZG dd„ deƒZ	dS )é    Né   )ÚPubSubManagerÚsocketioc                   sR   e Zd ZdZdZd‡ fdd„	Z‡ fd	d
„Zdd„ Zdd„ Zdd„ Z	dd„ Z
‡  ZS )ÚRedisManageraC  Redis based client manager.

    This class implements a Redis backend for event sharing across multiple
    processes. Only kept here as one more example of how to build a custom
    backend, since the kombu backend is perfectly adequate to support a Redis
    message queue.

    To use a Redis backend, initialize the :class:`Server` instance as
    follows::

        url = 'redis://hostname:port/0'
        server = socketio.Server(client_manager=socketio.RedisManager(url))

    :param url: The connection URL for the Redis server. For a default Redis
                store running on the same host, use ``redis://``.
    :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 ot ``True``, only initialize to emit events. The
                       default of ``False`` initializes the class for emitting
                       and receiving.
    :param redis_options: additional keyword arguments to be passed to
                          ``Redis.from_url()``.
    Úredisúredis://localhost:6379/0r   FNc                sB   t d krtdƒ‚|| _|pi | _| jƒ  tt| ƒj|||d� d S )NzLRedis package is not installed (Run "pip install redis" in your virtualenv).)ÚchannelÚ
write_onlyÚlogger)r   ÚRuntimeErrorÚ	redis_urlÚredis_optionsÚ_redis_connectÚsuperr   Ú__init__)ÚselfÚurlr   r	   r
   r   )Ú	__class__© úA/tmp/pip-build-mqc4i71p/python-socketio/socketio/redis_manager.pyr   )   s    
zRedisManager.__init__c                sl   t t| ƒjƒ  d}| jjdkr4ddlm} |dƒ}n d| jjkrTddlm} |dƒ}|sht	d| jj ƒ‚d S )	NTZeventletr   )Úis_monkey_patchedÚsocketZgevent)Úis_module_patchedz<Redis requires a monkey patched socket library to work with )
r   r   Ú
initializeÚserverZ
async_modeZeventlet.patcherr   Zgevent.monkeyr   r   )r   Zmonkey_patchedr   r   )r   r   r   r   6   s    
zRedisManager.initializec             C   s&   t jj| jf| jŽ| _ | j jƒ | _d S )N)r   ZRedisZfrom_urlr   r   Úpubsub)r   r   r   r   r   E   s    
zRedisManager._redis_connectc             C   sj   d}x`y"|s| j ƒ  | jj| jtj|ƒƒS  tjjk
r`   |rPtj	dƒ d}ntj	dƒ P Y qX qW d S )NTz#Cannot publish to redis... retryingFz$Cannot publish to redis... giving up)
r   r   Úpublishr   ÚpickleÚdumpsÚ
exceptionsÚConnectionErrorr
   Úerror)r   ÚdataÚretryr   r   r   Ú_publishJ   s    

zRedisManager._publishc             c   s–   d}d}xˆy8|r&| j ƒ  | jj| jƒ x| jjƒ D ]
}|V  q2W W q
 tjjk
rŒ   tj	dj
|ƒƒ d}tj|ƒ |d9 }|dkrˆd}Y q
X q
W d S )Nr   Fz0Cannot receive from redis... retrying in {} secsTé   é<   )r   r   Ú	subscriber   Úlistenr   r   r    r
   r!   ÚformatÚtimeÚsleep)r   Zretry_sleepÚconnectÚmessager   r   r   Ú_redis_listen_with_retriesY   s"    
z'RedisManager._redis_listen_with_retriesc             c   sh   | j jdƒ}| jj| j ƒ x:| jƒ D ].}|d |kr$|d dkr$d|kr$|d V  q$W | jj| j ƒ d S )Nzutf-8r   Útyper-   r"   )r   Úencoder   r'   r.   Zunsubscribe)r   r   r-   r   r   r   Ú_listenl   s    zRedisManager._listen)r   r   FNN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Únamer   r   r   r$   r.   r1   Ú__classcell__r   r   )r   r   r      s    r   )
Úloggingr   r*   r   ÚImportErrorZpubsub_managerr   Ú	getLoggerr
   r   r   r   r   r   Ú<module>   s   

