3
ñ\ð]ú  ã               @   sj   d dl Z d dlZd dlmZ yd dlZW n ek
r@   dZY nX ddlmZ dd„ ZG dd„ deƒZ	dS )	é    N)Úurlparseé   )ÚAsyncPubSubManagerc             C   sj   t | ƒ}|jdkrtdƒ‚|jdk}|jp,d}|jp6d}|j}|jrXt|jdd … ƒ}nd}|||||fS )	NÚredisÚredisszInvalid redis urlÚ	localhostië  r   r   >   r   r   )r   ÚschemeÚ
ValueErrorÚhostnameÚportÚpasswordÚpathÚint)ÚurlÚpÚsslÚhostr   r   Údb© r   úI/tmp/pip-build-mqc4i71p/python-socketio/socketio/asyncio_redis_manager.pyÚ_parse_redis_url   s    



r   c                   s6   e Zd ZdZdZd‡ fdd„	Zd	d
„ Zdd„ Z‡  ZS )ÚAsyncRedisManagera  Redis based client manager for asyncio servers.

    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::

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

    :param url: The connection URL for the Redis server. For a default Redis
                store running on the same host, use ``redis://``.  To use an
                SSL connection, use ``rediss://``.
    :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.
    Úaioredisúredis://localhost:6379/0ÚsocketioFNc                sN   t d krtdƒ‚t|ƒ\| _| _| _| _| _d | _d | _	t
ƒ j|||d� d S )NzORedis package is not installed (Run "pip install aioredis" in your virtualenv).)ÚchannelÚ
write_onlyÚlogger)r   ÚRuntimeErrorr   r   r   r   r   r   ÚpubÚsubÚsuperÚ__init__)Úselfr   r   r   r   )Ú	__class__r   r   r"   5   s    zAsyncRedisManager.__init__c             Ã   s¦   d}xœyN| j d kr:tj| j| jf| j| j| jd�I d H | _ | j j| j	t
j|ƒƒI d H S  tjtfk
rœ   |rˆ| jƒ jdƒ d | _ d}n| jƒ jdƒ P Y qX qW d S )NT)r   r   r   z#Cannot publish to redis... retryingFz$Cannot publish to redis... giving up)r   r   Úcreate_redisr   r   r   r   r   Úpublishr   ÚpickleÚdumpsÚ
RedisErrorÚOSErrorÚ_get_loggerÚerror)r#   ÚdataÚretryr   r   r   Ú_publishB   s     

zAsyncRedisManager._publishc             Ã   sÄ   d}xºy\| j d kr:tj| j| jf| j| j| jd�I d H | _ | j j| j	ƒI d H d | _
| j
jƒ I d H S  tjtfk
rº   | jƒ jdj|ƒƒ d | _ tj|ƒI d H  |d9 }|dkr¶d}Y qX qW d S )Nr   )r   r   r   r   z0Cannot receive from redis... retrying in {} secsé   é<   )r    r   r%   r   r   r   r   r   Ú	subscriber   ÚchÚgetr)   r*   r+   r,   ÚformatÚasyncioÚsleep)r#   Zretry_sleepr   r   r   Ú_listenX   s"    
zAsyncRedisManager._listen)r   r   FN)	Ú__name__Ú
__module__Ú__qualname__Ú__doc__Únamer"   r/   r8   Ú__classcell__r   r   )r$   r   r      s    r   )
r6   r'   Úurllib.parser   r   ÚImportErrorZasyncio_pubsub_managerr   r   r   r   r   r   r   Ú<module>   s   
