3
ñ\ð]±  ã               @   sL   d dl mZ d dlZd dlZd dlZd dlZddlmZ G dd„ deƒZdS )é    )ÚpartialNé   )ÚBaseManagerc                   sŠ   e Zd ZdZdZd‡ fdd„	Z‡ fdd	„Zd‡ fd
d„	Zddd„Zdd„ Z	dd„ Z
‡ fdd„Zdd„ Zdd„ Z‡ fdd„Zdd„ Z‡  ZS )ÚPubSubManagera=  Manage a client list attached to a pub/sub backend.

    This is a base class that enables multiple servers to share the list of
    clients, with the servers communicating events through a pub/sub backend.
    The use of a pub/sub backend also allows any client connected to the
    backend to emit events addressed to Socket.IO clients.

    The actual backends must be implemented by subclasses, this class only
    provides a pub/sub generic framework.

    :param channel: The channel name on which the server sends and receives
                    notifications.
    ZpubsubÚsocketioFNc                s0   t t| ƒjƒ  || _|| _tjƒ j| _|| _	d S )N)
Úsuperr   Ú__init__ÚchannelÚ
write_onlyÚuuidÚuuid4ÚhexÚhost_idÚlogger)Úselfr	   r
   r   )Ú	__class__© úB/tmp/pip-build-mqc4i71p/python-socketio/socketio/pubsub_manager.pyr      s
    zPubSubManager.__init__c                s<   t t| ƒjƒ  | js$| jj| jƒ| _| jƒ j	| j
d ƒ d S )Nz backend initialized.)r   r   Ú
initializer
   ÚserverZstart_background_taskÚ_threadÚthreadZ_get_loggerÚinfoÚname)r   )r   r   r   r   "   s    zPubSubManager.initializec       	   
      s˜   |j dƒr&tt| ƒj||||||d�S |p,d}|dk	rr| jdkrHtdƒ‚|dkrXtdƒ‚| j|||ƒ}|||f}nd}| jd||||||| j	dœƒ dS )	a/  Emit a message to a single client, a room, or all the clients
        connected to the namespace.

        This method takes care or propagating the message to all the servers
        that are connected through the message queue.

        The parameters are the same as in :meth:`.Server.emit`.
        Zignore_queue)Ú	namespaceÚroomÚskip_sidÚcallbackú/Nz:Callbacks can only be issued from the context of a server.z'Cannot use callback without a room set.Úemit)ÚmethodÚeventÚdatar   r   r   r   r   )
Úgetr   r   r   r   ÚRuntimeErrorÚ
ValueErrorZ_generate_ack_idÚ_publishr   )	r   r!   r"   r   r   r   r   ÚkwargsÚid)r   r   r   r   (   s"    





zPubSubManager.emitc             C   s   | j d||pddœƒ d S )NÚ
close_roomr   )r    r   r   )r&   )r   r   r   r   r   r   r)   F   s    zPubSubManager.close_roomc             C   s   t dƒ‚dS )z¤Publish a message on the Socket.IO channel.

        This method needs to be implemented by the different subclasses that
        support pub/sub backends.
        z.This method must be implemented in a subclass.N)ÚNotImplementedError)r   r"   r   r   r   r&   J   s    zPubSubManager._publishc             C   s   t dƒ‚dS )zãReturn the next message published on the Socket.IO channel,
        blocking until a message is available.

        This method needs to be implemented by the different subclasses that
        support pub/sub backends.
        z.This method must be implemented in a subclass.N)r*   )r   r   r   r   Ú_listenS   s    zPubSubManager._listenc                sz   |j dƒ}|j dƒ}|d k	r<t|ƒdkr<t| j|f|žŽ }nd }tt| ƒj|d |d |j dƒ|j dƒ|j dƒ|d	� d S )
Nr   r   é   r!   r"   r   r   r   )r   r   r   r   )r#   Úlenr   Ú_return_callbackr   r   r   )r   ÚmessageZremote_callbackZremote_host_idr   )r   r   r   Ú_handle_emit]   s    



zPubSubManager._handle_emitc             C   s^   | j |jdƒkrZy$|d }|d }|d }|d }W n tk
rH   d S X | j||||ƒ d S )Nr   Úsidr   r(   Úargs)r   r#   ÚKeyErrorZtrigger_callback)r   r/   r1   r   r(   r2   r   r   r   Ú_handle_callbackn   s    zPubSubManager._handle_callbackc             G   s   | j d|||||dœƒ d S )Nr   )r    r   r1   r   r(   r2   )r&   )r   r   r1   r   Zcallback_idr2   r   r   r   r.   y   s    zPubSubManager._return_callbackc                s$   t t| ƒj|jdƒ|jdƒd� d S )Nr   r   )r   r   )r   r   r)   r#   )r   r/   )r   r   r   Ú_handle_close_room€   s    
z PubSubManager._handle_close_roomc             C   sÈ   xÂ| j ƒ D ]¶}d }t|tƒr"|}nLt|tjƒrJytj|ƒ}W n   Y nX |d krnytj|ƒ}W n   Y nX |r
d|kr
|d dkr’| j|ƒ q
|d dkrª| j	|ƒ q
|d dkr
| j
|ƒ q
W d S )Nr    r   r   r)   )r+   Ú
isinstanceÚdictÚsixÚbinary_typeÚpickleÚloadsÚjsonr0   r4   r5   )r   r/   r"   r   r   r   r   „   s*    
zPubSubManager._thread)r   FN)NNNN)N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r   r   r)   r&   r+   r0   r4   r.   r5   r   Ú__classcell__r   r   )r   r   r      s    
	
r   )	Ú	functoolsr   r   r<   r:   r8   Zbase_managerr   r   r   r   r   r   Ú<module>   s   