3
ò\ð]•i  ã               @   s  d dl Z yd dlZW n ek
r0   d dlZY nX d dlZd dlZd dlZ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 yd dlZW n ek
r°   dZY nX ddlmZ ddlmZ ddlmZ e jdƒZg ZejrîeZdd	„ ZejejeƒZG d
d„ deƒZdS )é    N)Úurllibé   )Ú
exceptions)Úpacket)Úpayloadzengineio.clientc             C   s^   x:t dd… D ]*}|jƒ r,|j|jdd� q|jdd� qW ttƒrNt| |ƒS tj| |ƒS dS )zdSIGINT handler.

    Disconnect all active clients and then invoke the original signal handler.
    NT)Úabort)Úconnected_clientsÚis_asyncio_basedÚstart_background_taskÚ
disconnectÚcallableÚoriginal_signal_handlerÚsignalÚdefault_int_handler)ÚsigÚframeÚclient© r   ú:/tmp/pip-build-mqc4i71p/python-engineio/engineio/client.pyÚsignal_handler    s    
r   c               @   sö   e Zd ZdZdddgZd=d	d
„Zdd„ Zd>dd„Zi ddfdd„Zdd„ Z	d?dd„Z
d@dd„Zdd„ Zdd„ ZdAdd„Zdd „ Zd!d"„ Zd#d$„ Zd%d&„ Zd'd(„ Zd)d*„ Zd+d,„ ZdBd-d.„Zd/d0„ Zd1d2„ Zd3d4„ Zd5d6„ Zd7d8„ Zd9d:„ Zd;d<„ ZdS )CÚClientaÔ  An Engine.IO client.

    This class implements a fully compliant Engine.IO web client with support
    for websocket and long-polling transports.

    :param logger: To enable logging set to ``True`` or pass a logger object to
                   use. To disable logging set to ``False``. The default is
                   ``False``.
    :param json: An alternative json module to use for encoding and decoding
                 packets. Custom json modules must have ``dumps`` and ``loads``
                 functions that are compatible with the standard library
                 versions.
    :param request_timeout: A timeout in seconds for requests. The default is
                            5 seconds.
    :param ssl_verify: ``True`` to verify SSL certificates, or ``False`` to
                       skip SSL certificate verification, allowing
                       connections to servers with self signed certificates.
                       The default is ``True``.
    Úconnectr   ÚmessageFNé   Tc             C   sè   i | _ d | _d | _d | _d | _d | _d | _d | _d| _d | _	d | _
d | _d | _d | _d | _d | _d| _|| _|d k	r||tj_t|tƒsŽ|| _nPt| _tjj  rÞ| jjtjkrÞ|rÀ| jjtjƒ n| jjtjƒ | jj tj!ƒ ƒ || _"d S )NTÚdisconnected)#ÚhandlersÚbase_urlÚ
transportsÚcurrent_transportÚsidÚupgradesÚping_intervalÚping_timeoutÚpong_receivedÚhttpÚwsÚread_loop_taskÚwrite_loop_taskÚping_loop_taskÚping_loop_eventÚqueueÚstateÚ
ssl_verifyr   ÚPacketÚjsonÚ
isinstanceÚboolÚloggerÚdefault_loggerÚloggingÚrootÚlevelÚNOTSETÚsetLevelÚINFOÚERRORÚ
addHandlerÚStreamHandlerÚrequest_timeout)Úselfr1   r.   r<   r,   r   r   r   Ú__init__J   s<    

zClient.__init__c             C   s   dS )NFr   )r=   r   r   r   r	   r   s    zClient.is_asyncio_basedc                s8   ˆ ˆj krtdƒ‚‡ ‡fdd„}|dkr,|S ||ƒ dS )aâ  Register an event handler.

        :param event: The event name. Can be ``'connect'``, ``'message'`` or
                      ``'disconnect'``.
        :param handler: The function that should be invoked to handle the
                        event. When this parameter is not given, the method
                        acts as a decorator for the handler function.

        Example usage::

            # as a decorator:
            @eio.on('connect')
            def connect_handler():
                print('Connection request')

            # as a method:
            def message_handler(msg):
                print('Received message: ', msg)
                eio.send('response')
            eio.on('message', message_handler)
        zInvalid eventc                s   | ˆj ˆ < | S )N)r   )Úhandler)Úeventr=   r   r   Úset_handlerŽ   s    
zClient.on.<locals>.set_handlerN)Úevent_namesÚ
ValueError)r=   r@   r?   rA   r   )r@   r=   r   Úonu   s    
z	Client.onz	engine.ioc                s‚   | j dkrtdƒ‚ddg‰ |dk	rRt|tjƒr4|g}‡ fdd„|D ƒ}|sRtdƒ‚|pXˆ | _| jƒ | _t| d	| jd
  ƒ|||ƒS )aˆ  Connect to an Engine.IO server.

        :param url: The URL of the Engine.IO server. It can include custom
                    query string parameters if required by the server.
        :param headers: A dictionary with custom headers to send with the
                        connection request.
        :param transports: The list of allowed transports. Valid transports
                           are ``'polling'`` and ``'websocket'``. If not
                           given, the polling transport is connected first,
                           then an upgrade to websocket is attempted.
        :param engineio_path: The endpoint where the Engine.IO server is
                              installed. The default value is appropriate for
                              most cases.

        Example usage::

            eio = engineio.Client()
            eio.connect('http://localhost:5000')
        r   z%Client is not in a disconnected stateÚpollingÚ	websocketNc                s   g | ]}|ˆ kr|‘qS r   r   )Ú.0Ú	transport)Úvalid_transportsr   r   ú
<listcomp>±   s    z"Client.connect.<locals>.<listcomp>zNo valid transports providedZ	_connect_r   )	r+   rC   r/   ÚsixÚstring_typesr   Úcreate_queuer*   Úgetattr)r=   ÚurlÚheadersr   Úengineio_pathr   )rI   r   r   –   s    


zClient.connectc             C   s   | j r| j jƒ  dS )z¯Wait until the connection with the server ends.

        Client applications can use this function to block the main thread
        during the life of the connection.
        N)r&   Újoin)r=   r   r   r   Úwaitº   s    zClient.waitc             C   s   | j tjtj||d�ƒ dS )a  Send a message to a client.

        :param data: The data to send to the client. Data can be of type
                     ``str``, ``bytes``, ``list`` or ``dict``. If a ``list``
                     or ``dict``, the data will be serialized as JSON.
        :param binary: ``True`` to send packet as binary, ``False`` to send
                       as text. If not given, unicode (Python 2) and str
                       (Python 3) are sent as text, and str (Python 2) and
                       bytes (Python 3) are sent as binary.
        )ÚdataÚbinaryN)Ú_send_packetr   r-   ÚMESSAGE)r=   rT   rU   r   r   r   ÚsendÃ   s    zClient.sendc             C   s”   | j dkrˆ| jtjtjƒƒ | jjdƒ d| _ | jddd� | jdkrP| j	j
ƒ  |s^| jjƒ  d| _ ytj| ƒ W n tk
r†   Y nX | jƒ  dS )	z­Disconnect from the server.

        :param abort: If set to ``True``, do not wait for background tasks
                      associated with the connection to end.
        Ú	connectedNZdisconnectingr   F)Ú	run_asyncrF   r   )r+   rV   r   r-   ÚCLOSEr*   ÚputÚ_trigger_eventr   r%   Úcloser&   rR   r   ÚremoverC   Ú_reset)r=   r   r   r   r   r   Ñ   s    



zClient.disconnectc             C   s   | j S )z¡Return the name of the transport currently in use.

        The possible values returned by this function are ``'polling'`` and
        ``'websocket'``.
        )r   )r=   r   r   r   rH   ç   s    zClient.transportc             O   s   t j|||d�}|jƒ  |S )aù  Start a background task.

        This is a utility function that applications can use to start a
        background task.

        :param target: the target function to execute.
        :param args: arguments to pass to the function.
        :param kwargs: keyword arguments to pass to the function.

        This function returns an object compatible with the `Thread` class in
        the Python standard library. The `start()` method on this object is
        already called by this function.
        )ÚtargetÚargsÚkwargs)Ú	threadingÚThreadÚstart)r=   ra   rb   rc   Úthr   r   r   r
   ï   s    zClient.start_background_taskr   c             C   s
   t j|ƒS )z'Sleep for the requested amount of time.)ÚtimeÚsleep)r=   Úsecondsr   r   r   ri     s    zClient.sleepc             O   s   t j||Ž}t j|_|S )zCreate a queue object.)r*   ÚQueueÚEmpty)r=   rb   rc   Úqr   r   r   rM     s    zClient.create_queuec             O   s   t j||ŽS )zCreate an event object.)rd   ÚEvent)r=   rb   rc   r   r   r   Úcreate_event  s    zClient.create_eventc             C   s   d| _ d | _d S )Nr   )r+   r   )r=   r   r   r   r`     s    zClient._resetc             C   sö  t dkr| jjdƒ dS | j||dƒ| _| jjd| j ƒ | jd| j| jƒ  || jd�}|dkrr| j	ƒ  t
jdƒ‚|jdk s†|jd	kr˜t
jd
j|jƒƒ‚ytj|jd�}W n& tk
rÐ   tjt
jdƒdƒ Y nX |jd }|jtjkròt
jdƒ‚| jjdt|jƒ ƒ |jd | _|jd | _|jd d | _|jd d | _d| _|  jd| j 7  _d| _t j!| ƒ | j"ddd� x"|jdd… D ]}| j#|ƒ �qˆW d| jk�rÈd| j$k�rÈ| j%|||ƒ�rÈdS | j&| j'ƒ| _(| j&| j)ƒ| _*| j&| j+ƒ| _,dS )z<Establish a long-polling connection to the Engine.IO server.Nz?requests package is not installed -- cannot send HTTP requests!rE   z!Attempting polling connection to ÚGET)rP   Útimeoutz Connection refused by the serveréÈ   i,  z,Unexpected status code {} in server response)Úencoded_payloadzUnexpected response from serverr   z"OPEN packet not returned by serverz!Polling connection accepted with r   r    ÚpingIntervalg     @�@ÚpingTimeoutz&sid=rY   r   F)rZ   r   rF   )-Úrequestsr1   ÚerrorÚ_get_engineio_urlr   ÚinfoÚ_send_requestÚ_get_url_timestampr<   r`   r   ÚConnectionErrorÚstatus_codeÚformatr   ÚPayloadÚcontentrC   rK   Ú
raise_fromÚpacketsÚpacket_typer   ÚOPENÚstrrT   r   r    r!   r"   r   r+   r   Úappendr]   Ú_receive_packetr   Ú_connect_websocketr
   Ú
_ping_loopr(   Ú_write_loopr'   Ú_read_loop_pollingr&   )r=   rO   rP   rQ   ÚrÚpÚopen_packetÚpktr   r   r   Ú_connect_polling  sZ    



zClient._connect_pollingc          4   C   s\  t dkr| jjdƒ dS | j||dƒ}| jrP| jjd| ƒ d}|d| j 7 }nd}|| _| jjd| ƒ d}| jrŒd	jd
d„ | jj	D ƒƒ}yD| j
s¶t j|| jƒ  ||dtjid�}nt j|| jƒ  ||d�}W n8 ttfk
�r   |rú| jjdƒ dS tjdƒ‚Y nX |�rNtjtjtjdƒd�jƒ }y|j|ƒ W n4 tk
�rl }	 z| jjdt|	ƒƒ dS d}	~	X nX y|jƒ }W n4 tk
�r® }	 z| jjdt|	ƒƒ dS d}	~	X nX tj|d�}
|
jtjk�sÖ|
jdk�ræ| jjdƒ dS tjtjƒjƒ }y|j|ƒ W n4 tk
�r8 }	 z| jjdt|	ƒƒ dS d}	~	X nX d| _ | jjdƒ nÚy|jƒ }W n6 tk
�r� }	 ztjdt|	ƒ ƒ‚W Y dd}	~	X nX tj|d�}|jtj!k�r¶tjdƒ‚| jjdt|jƒ ƒ |jd | _|jd | _"|jd d | _#|jd d | _$d| _ d | _%t&j'| ƒ | j(d!dd"� || _)| j*| j+ƒ| _,| j*| j-ƒ| _.| j*| j/ƒ| _0dS )#z?Establish or upgrade to a WebSocket connection with the server.NzKwebsocket-client package not installed, only polling transport is availableFrF   z Attempting WebSocket upgrade to Tz&sid=z#Attempting WebSocket connection to z; c             S   s   g | ]}d j |j|jƒ‘qS )z{}={})r~   ÚnameÚvalue)rG   Úcookier   r   r   rJ   c  s   z-Client._connect_websocket.<locals>.<listcomp>Ú	cert_reqs)Úheaderr“   Zsslopt)r•   r“   z*WebSocket upgrade failed: connection errorzConnection errorZprobe)rT   z7WebSocket upgrade failed: unexpected send exception: %sz7WebSocket upgrade failed: unexpected recv exception: %s)Úencoded_packetz(WebSocket upgrade failed: no PONG packetz WebSocket upgrade was successfulzUnexpected recv exception: zno OPEN packetz#WebSocket connection accepted with r   r    rt   g     @�@ru   rY   r   )rZ   )1rF   r1   Úwarningrx   r   ry   r   r$   rR   Úcookiesr,   Úcreate_connectionr{   ÚsslÚ	CERT_NONEr|   ÚIOErrorr   r   r-   ÚPINGrK   Ú	text_typeÚencoderX   Ú	Exceptionr…   Úrecvrƒ   ÚPONGrT   ÚUPGRADEr   r„   r    r!   r"   r+   r   r†   r]   r%   r
   r‰   r(   rŠ   r'   Ú_read_loop_websocketr&   )r=   rO   rP   rQ   Zwebsocket_urlÚupgrader˜   r%   r�   Úer�   rŽ   r   r   r   rˆ   L  s®    





 


zClient._connect_websocketc             C   s²   |j ttjƒk rtj|j  nd}| jjd|t|jtƒs<|jndƒ |j tj	krb| j
d|jdd� nL|j tjkrvd| _n8|j tjkr�| jdd� n|j tjkržn| jjd|j ƒ d	S )
z(Handle incoming packets from the server.ÚUNKNOWNzReceived packet %s data %sz<binary>r   T)rZ   )r   z%Received unexpected packet of type %sN)rƒ   Úlenr   Úpacket_namesr1   ry   r/   rT   ÚbytesrW   r]   r¢   r#   r[   r   ZNOOPrw   )r=   r�   Zpacket_namer   r   r   r‡   ³  s     zClient._receive_packetc             C   sH   | j dkrdS | jj|ƒ | jjdtj|j t|j	t
ƒs>|j	ndƒ dS )z(Queue a packet to be sent to the server.rY   NzSending packet %s data %sz<binary>)r+   r*   r\   r1   ry   r   r©   rƒ   r/   rT   rª   )r=   r�   r   r   r   rV   Æ  s    

zClient._send_packetc             C   sl   | j d krtjƒ | _ y| j j|||||| jd�S  tjjk
rf } z| jjd|||ƒ W Y d d }~X nX d S )N)rP   rT   rq   Úverifyz+HTTP %s request to %s failed with error %s.)	r$   rv   ÚSessionÚrequestr,   r   ÚRequestExceptionr1   ry   )r=   ÚmethodrO   rP   Úbodyrq   Úexcr   r   r   rz   Ð  s    

zClient._send_requestc          	   O   s`   |j ddƒ}|| jkr\|r0| j| j| f|žŽ S y| j| |Ž S    | jj|d ƒ Y nX dS )zInvoke an event handler.rZ   Fz handler errorN)Úpopr   r
   r1   Ú	exception)r=   r@   rb   rc   rZ   r   r   r   r]   Ü  s    
zClient._trigger_eventc             C   sp   |j dƒ}tjj|ƒ}|dkr$d}n|dkr2d}ntdƒ‚|jdkrL|d	7 }d
j||j||j|jrfdnd|d�S )z&Generate the Engine.IO connection URL.ú/rE   r$   rF   r%   zinvalid transportÚhttpsÚwssÚszC{scheme}://{netloc}/{path}/?{query}{sep}transport={transport}&EIO=3ú&Ú )ÚschemeÚnetlocÚpathÚqueryÚseprH   )rµ   r¶   )	Ústripr   ÚparseÚurlparserC   rº   r~   r»   r½   )r=   rO   rQ   rH   Ú
parsed_urlrº   r   r   r   rx   è  s    

zClient._get_engineio_urlc             C   s   dt tjƒ ƒ S )z.Generate the Engine.IO query string timestamp.z&t=)r…   rh   )r=   r   r   r   r{   ý  s    zClient._get_url_timestampc             C   sž   d| _ | jdkr| jƒ | _n
| jjƒ  xf| jdkrŒ| j sb| jjdƒ | jrT| jjƒ  | j	j
dƒ P d| _ | jtjtjƒƒ | jj| jd� q(W | jjdƒ dS )z[This background task sends a PING to the server at the requested
        interval.
        TNrY   z-PONG response has not been received, abortingF)rq   zExiting ping task)r#   r)   ro   Úclearr+   r1   ry   r%   Úshutdownr*   r\   rV   r   r-   r�   rS   r!   )r=   r   r   r   r‰     s     


zClient._ping_loopc             C   s�  xø| j dkrø| jjd| j ƒ | jd| j| jƒ  t| j| jƒd d�}|dkrh| jj	dƒ | j
jdƒ P |jdk s||jd	krš| jj	d
|jƒ | j
jdƒ P ytj|jd�}W n. tk
rÚ   | jj	dƒ | j
jdƒ P Y nX x|jD ]}| j|ƒ qäW qW | jjdƒ | jjƒ  | jjdƒ | j�r.| jjƒ  | jjƒ  | j dk�r€| jddd� ytj| ƒ W n tk
�rv   Y nX | jƒ  | jjdƒ dS )z-Read packets by polling the Engine.IO server.rY   zSending polling GET request to rp   r   )rq   Nz*Connection refused by the server, abortingrr   i,  z6Unexpected status code %s in server response, aborting)rs   z'Unexpected packet from server, abortingz"Waiting for write loop task to endz!Waiting for ping loop task to endr   F)rZ   zExiting read loop task)r+   r1   ry   r   rz   r{   Úmaxr!   r"   r—   r*   r\   r}   r   r   r€   rC   r‚   r‡   r'   rR   r)   Úsetr(   r]   r   r_   r`   )r=   rŒ   r�   r�   r   r   r   r‹     sN    


zClient._read_loop_pollingc             C   sT  x¾| j dkr¾d}y| jjƒ }W np tjk
rN   | jjdƒ | jjdƒ P Y nB t	k
rŽ } z&| jj
dt|ƒƒ | jjdƒ P W Y dd}~X nX t|tjƒr¦|jdƒ}tj|d�}| j|ƒ qW | jj
dƒ | jjƒ  | jj
dƒ | jrò| jjƒ  | jjƒ  | j dk�rD| jd	d
d� ytj| ƒ W n tk
�r:   Y nX | jƒ  | jj
dƒ dS )z5Read packets from the Engine.IO WebSocket connection.rY   Nz)WebSocket connection was closed, abortingzUnexpected error "%s", abortingzutf-8)r–   z"Waiting for write loop task to endz!Waiting for ping loop task to endr   F)rZ   zExiting read loop task)r+   r%   r¡   rF   Ú"WebSocketConnectionClosedExceptionr1   r—   r*   r\   r    ry   r…   r/   rK   rž   rŸ   r   r-   r‡   r'   rR   r)   rÆ   r(   r]   r   r_   rC   r`   )r=   r�   r¦   r�   r   r   r   r¤   B  s@    



zClient._read_loop_websocketc             C   s  �xò| j dk�rôt| j| jƒd }d}y| jj|d�g}W n& | jjk
r`   | jjdƒ P Y nX |dgkr|| jj	ƒ  g }n^x\y|j
| jjdd�ƒ W n | jjk
r°   P Y nX |d dkr~|dd… }| jj	ƒ  P q~W |sàP | jd	k�r~tj|d
�}| jd| j|jƒ ddi| jd�}x|D ]}| jj	ƒ  �qW |dk�rJ| jjdƒ P |jdk �sb|jdk�rò| jjd|jƒ | jƒ  P qyLxF|D ]>}|jdd�}|j�r¬| jj|ƒ n| jj|ƒ | jj	ƒ  �q†W W q tjk
�rð   | jjdƒ P Y qX qW | jjdƒ dS )zhThis background task sends packages to the server as they are
        pushed to the send queue.
        rY   r   N)rq   zpacket queue is empty, abortingF)Úblockr   rE   )r‚   ÚPOSTzContent-Typezapplication/octet-stream)r°   rP   rq   z*Connection refused by the server, abortingrr   i,  z6Unexpected status code %s in server response, aborting)Zalways_bytesz)WebSocket connection was closed, abortingzExiting write loop taskéÿÿÿÿrÊ   )r+   rÅ   r!   r"   r*   Úgetrl   r1   rw   Ú	task_doner†   r   r   r   rz   r   rŸ   r<   r—   r}   r`   rU   r%   Zsend_binaryrX   rF   rÇ   ry   )r=   rq   r‚   r�   rŒ   r�   r–   r   r   r   rŠ   f  sf    






zClient._write_loop)FNr   T)N)N)F)r   )NNN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__rB   r>   r	   rD   r   rS   rX   r   rH   r
   ri   rM   ro   r`   r�   rˆ   r‡   rV   rz   r]   rx   r{   r‰   r‹   r¤   rŠ   r   r   r   r   r   4   s@   
   
$
!#	


9g 

+$r   )r3   r*   ÚImportErrorrk   r   rš   rd   rh   rK   Z	six.movesr   rv   rF   r¹   r   r   r   Ú	getLoggerr2   r   ÚPY2ÚOSErrorr|   r   ÚSIGINTr   Úobjectr   r   r   r   r   Ú<module>   s8   


