3
ò\ð]Ê_  ã               @   s„   d dl Z d dlZyd dlZW n ek
r4   dZY nX d dlZddlmZ ddlmZ ddlmZ ddlm	Z	 G dd„ dej
ƒZdS )	é    Né   )Úclient)Ú
exceptions)Úpacket)Úpayloadc                   sÈ   e Zd ZdZdd„ Zi ddfdd„Zdd	„ Zd.d
d„Zd/dd„Zdd„ Z	d0dd„Z
dd„ Zdd„ Z‡ fdd„Zdd„ Zdd„ Zdd„ Zd d!„ Zd1d"d#„Zd$d%„ Zd&d'„ Zd(d)„ Zd*d+„ Zd,d-„ Z‡  ZS )2ÚAsyncClienta"  An Engine.IO client for asyncio.

    This class implements a fully compliant Engine.IO web client with support
    for websocket and long-polling transports, compatible with the asyncio
    framework on Python 3.5 or newer.

    :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``.
    c             C   s   dS )NT© )Úselfr   r   úB/tmp/pip-build-mqc4i71p/python-engineio/engineio/asyncio_client.pyÚis_asyncio_based%   s    zAsyncClient.is_asyncio_basedNz	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
  ƒ|||ƒI dH 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.

        Note: this method is a coroutine.

        Example usage::

            eio = engineio.Client()
            await eio.connect('http://localhost:5000')
        Údisconnectedz%Client is not in a disconnected stateÚpollingÚ	websocketNc                s   g | ]}|ˆ kr|‘qS r   r   )Ú.0Ú	transport)Úvalid_transportsr   r
   ú
<listcomp>E   s    z'AsyncClient.connect.<locals>.<listcomp>zNo valid transports providedZ	_connect_r   )	ÚstateÚ
ValueErrorÚ
isinstanceÚsixÚ	text_typeÚ
transportsÚcreate_queueÚqueueÚgetattr)r	   ÚurlÚheadersr   Úengineio_pathr   )r   r
   Úconnect(   s    


zAsyncClient.connectc             Ã   s   | j r| j I dH  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.

        Note: this method is a coroutine.
        N)Úread_loop_task)r	   r   r   r
   ÚwaitN   s    zAsyncClient.waitc             Ã   s"   | j tjtj||d�ƒI dH  dS )aI  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.

        Note: this method is a coroutine.
        )ÚdataÚbinaryN)Ú_send_packetr   ÚPacketÚMESSAGE)r	   r"   r#   r   r   r
   ÚsendY   s    zAsyncClient.sendFc             Ã   s°   | j dkr¤| jtjtjƒƒI dH  | jjdƒI dH  d| _ | jddd�I dH  | jdkrh| j	j
ƒ I dH  |sx| jI dH  d| _ ytj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.

        Note: this method is a coroutine.
        Ú	connectedNZdisconnectingÚ
disconnectF)Ú	run_asyncr   r   )r   r$   r   r%   ÚCLOSEr   ÚputÚ_trigger_eventÚcurrent_transportÚwsÚcloser    r   Úconnected_clientsÚremover   Ú_reset)r	   Úabortr   r   r
   r)   i   s    

zAsyncClient.disconnectc             O   s   t 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.

        Note: this method is a coroutine.
        )ÚasyncioÚensure_future)r	   ÚtargetÚargsÚkwargsr   r   r
   Ústart_background_task�   s    z!AsyncClient.start_background_taskr   c             Ã   s   t j|ƒI dH S )z[Sleep for the requested amount of time.

        Note: this method is a coroutine.
        N)r5   Úsleep)r	   Úsecondsr   r   r
   r;   “   s    zAsyncClient.sleepc             C   s   t jƒ }t j|_|S )zCreate a queue object.)r5   ÚQueueZ
QueueEmptyÚEmpty)r	   Úqr   r   r
   r   š   s    zAsyncClient.create_queuec             C   s   t jƒ S )zCreate an event object.)r5   ÚEvent)r	   r   r   r
   Úcreate_event    s    zAsyncClient.create_eventc                s$   | j rtj| j jƒ ƒ tƒ jƒ  d S )N)Úhttpr5   r6   r0   Úsuperr3   )r	   )Ú	__class__r   r
   r3   ¤   s    zAsyncClient._resetc             Ã   s  t dkr| jjdƒ dS | j||dƒ| _| jjd| j ƒ | jd| j| jƒ  || jd�I dH }|dkrx| j	ƒ  t
jdƒ‚|jdk sŒ|jd	kržt
jd
j|jƒƒ‚ytj|jƒ I dH 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"| ƒ | j#ddd�I dH  x(|jdd… D ]}| j$|ƒI dH  �q W d| jk�rìd| j%k�rì| j&|||ƒI dH �rìdS | j'| j(ƒ| _)| j'| j*ƒ| _+| j'| j,ƒ| _-dS )z<Establish a long-polling connection to the Engine.IO server.Nz3aiohttp not installed -- cannot make HTTP requests!r   z!Attempting polling connection to ÚGET)r   Ú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 ÚsidÚupgradesÚpingIntervalg     @�@ÚpingTimeoutz&sid=r(   r   F)r*   r   r   ).ÚaiohttpÚloggerÚerrorÚ_get_engineio_urlÚbase_urlÚinfoÚ_send_requestÚ_get_url_timestampÚrequest_timeoutr3   r   ÚConnectionErrorÚstatusÚformatr   ÚPayloadÚreadr   r   Ú
raise_fromÚpacketsÚpacket_typer   ÚOPENÚstrr"   rI   rJ   Úping_intervalÚping_timeoutr.   r   r   r1   Úappendr-   Ú_receive_packetr   Ú_connect_websocketr:   Ú
_ping_loopÚping_loop_taskÚ_write_loopÚwrite_loop_taskÚ_read_loop_pollingr    )r	   r   r   r   ÚrÚpÚopen_packetÚpktr   r   r
   Ú_connect_polling©   sZ    

zAsyncClient._connect_pollingc          4   Ã   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| ƒ y`| jsªtj	ƒ }d|_
tj|_| jj|| jƒ  ||d	�I dH }n| jj|| jƒ  |d
�I dH }W nB t jjt jjfk
�r   |� rþ| jjdƒ dS tjdƒ‚Y nX |�rhtjtjdd�jdd�}y|j|ƒI dH  W n4 tk
�rt }	 z| jjdt|	ƒƒ dS d}	~	X nX y|jƒ I dH 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dd�}y|j|ƒI dH  W n4 tk
�rR }	 z| jjdt|	ƒƒ dS d}	~	X nX d| _"| jjdƒ nêy|jƒ I dH 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*| ƒ | j+ddd�I dH  || _,| j-| j.ƒ| _/| j-| j0ƒ| _1| j-| j2ƒ| _3dS ) z?Establish or upgrade to a WebSocket connection with the server.Nzaiohttp package not installedFr   z Attempting WebSocket upgrade to Tz&sid=z#Attempting WebSocket connection to )r   Ússl)r   z*WebSocket upgrade failed: connection errorzConnection errorZprobe)r"   )Úalways_bytesz7WebSocket 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 rI   rJ   rK   g     @�@rL   r(   r   )r*   )4rM   rN   rO   rP   rI   rR   rQ   Ú
ssl_verifyro   Úcreate_default_contextÚcheck_hostnameÚ	CERT_NONEÚverify_moderB   Z
ws_connectrT   Úclient_exceptionsZWSServerHandshakeErrorZServerConnectionErrorÚwarningr   rV   r   r%   ÚPINGÚencodeÚsend_strÚ	Exceptionr_   Úreceiver"   r]   ÚPONGÚUPGRADEr.   r^   rJ   r`   ra   r   r   r1   rb   r-   r/   r:   re   rf   rg   rh   Ú_read_loop_websocketr    )r	   r   r   r   Zwebsocket_urlÚupgradeÚssl_contextr/   rk   Úerm   rl   r   r   r
   rd   à   s°    






 

zAsyncClient._connect_websocketc             Ã   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rh| j
d|jdd�I dH  nR|j tjkr|d| _n>|j tjkrœ| jdd�I dH  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>ÚmessageT)r*   N)r4   z%Received unexpected packet of type %s)r]   Úlenr   Úpacket_namesrN   rR   r   r"   Úbytesr&   r-   r~   Úpong_receivedr+   r)   ZNOOPrO   )r	   rm   Zpacket_namer   r   r
   rc   B  s     zAsyncClient._receive_packetc             Ã   sN   | j dkrdS | jj|ƒI dH  | jjdtj|j t|j	t
ƒsD|j	ndƒ dS )z(Queue a packet to be sent to the server.r(   NzSending packet %s data %sz<binary>)r   r   r,   rN   rR   r   r‡   r]   r   r"   rˆ   )r	   rm   r   r   r
   r$   U  s    

zAsyncClient._send_packetc             Ã   s¶   | j d ks| j jrtjƒ | _ t| j |jƒ ƒ}yH| jsT||||tj|d�dd�I d H S ||||tj|d�d�I d H S W n< tjt	j
fk
r° } z| jjd|||ƒ W Y d d }~X nX d S )N)ÚtotalF)r   r"   rF   ro   )r   r"   rF   z+HTTP %s request to %s failed with error %s.)rB   ÚclosedrM   ZClientSessionr   Úlowerrr   ZClientTimeoutZClientErrorr5   ÚTimeoutErrorrN   rR   )r	   Úmethodr   r   ÚbodyrF   Zhttp_methodÚexcr   r   r
   rS   _  s    
zAsyncClient._send_requestc             �   sþ   |j ddƒ}d}ˆˆjkrútjˆjˆ ƒdkr |rHˆjˆjˆ fˆ žŽ S yˆjˆ ˆ Ž I dH }W qú tjk
rv   Y qú   ˆjjˆd ƒ ˆdkr˜dS Y qúX nZ|r¾‡ ‡‡fdd„}ˆj|ƒS yˆjˆ ˆ Ž }W n(   ˆjjˆd	 ƒ ˆdkrôdS Y nX |S )
zInvoke an event handler.r*   FNTz async handler errorr   c               “   s   ˆj ˆ ˆ Ž S )N)Úhandlersr   )r8   Úeventr	   r   r
   Úasync_handlerŠ  s    z1AsyncClient._trigger_event.<locals>.async_handlerz handler error)Úpopr‘   r5   Úiscoroutinefunctionr:   ÚCancelledErrorrN   Ú	exception)r	   r’   r8   r9   r*   Úretr“   r   )r8   r’   r	   r
   r-   t  s2    


zAsyncClient._trigger_eventc             Ã   sÜ   d| _ | jdkr| jƒ | _n
| jjƒ  x¤| jdkrÊ| j sn| jjdƒ | jrZ| jjƒ I dH  | j	j
dƒI dH  P d| _ | jtjtjƒƒI dH  ytj| jjƒ | jƒI dH  W q( tjtjfk
rÆ   Y q(X q(W | jjdƒ dS )z[This background task sends a PING to the server at the requested
        interval.
        TNr(   z-PONG response has not been received, abortingFzExiting ping task)r‰   Úping_loop_eventrA   Úclearr   rN   rR   r/   r0   r   r,   r$   r   r%   ry   r5   Úwait_forr!   r`   r�   r–   )r	   r   r   r
   re   ™  s*    


zAsyncClient._ping_loopc             Ã   sÈ  �x"| j dk�r$| jjd| j ƒ | jd| j| jƒ  t| j| jƒd d�I dH }|dkrx| jj	dƒ | j
jdƒI dH  P |jdk sŒ|jd	kr°| jj	d
|jƒ | j
jdƒI dH  P ytj|jƒ I dH d�}W n4 tk
rþ   | jj	dƒ | j
jdƒI dH  P Y nX x |jD ]}| j|ƒI dH  �qW qW | jjdƒ | jI dH  | jjdƒ | j�r\| jjƒ  | jI dH  | j dk�r¸| jddd�I dH  ytjj| ƒ W n tk
�r®   Y nX | jƒ  | jjdƒ dS )z-Read packets by polling the Engine.IO server.r(   zSending polling GET request to rE   é   )rF   Nz*Connection refused by the server, abortingrG   i,  z6Unexpected status code %s in server response, aborting)rH   z'Unexpected packet from server, abortingz"Waiting for write loop task to endz!Waiting for ping loop task to endr)   F)r*   zExiting read loop task)r   rN   rR   rQ   rS   rT   Úmaxr`   ra   rx   r   r,   rW   r   rY   rZ   r   r\   rc   rh   r™   Úsetrf   r-   r   r1   r2   r3   )r	   rj   rk   rm   r   r   r
   ri   ´  sN    
zAsyncClient._read_loop_pollingc             Ã   s~  xÚ| j dkrÚd}y| jjƒ I dH j}W n~ tjjk
r^   | jjdƒ | j	j
dƒI dH  P Y nH tk
r¤ } z,| jjdt|ƒƒ | j	j
dƒI dH  P W Y dd}~X nX t|tjƒr¼|jdƒ}tj|d�}| j|ƒI dH  qW | jjdƒ | jI dH  | jjdƒ | j�r| jjƒ  | jI dH  | j dk�rn| jd	d
d�I dH  ytjj| ƒ W n tk
�rd   Y nX | jƒ  | jjdƒ dS )z5Read packets from the Engine.IO WebSocket connection.r(   Nz4Read loop: WebSocket connection was closed, abortingzUnexpected error "%s", abortingzutf-8)rq   z"Waiting for write loop task to endz!Waiting for ping loop task to endr)   F)r*   zExiting read loop task)r   r/   r}   r"   rM   rw   ÚServerDisconnectedErrorrN   rR   r   r,   r|   r_   r   r   r   rz   r   r%   rc   rh   r™   rž   rf   r-   r   r1   r2   r   r3   )r	   rk   rƒ   rm   r   r   r
   r€   ß  s@    

z AsyncClient._read_loop_websocketc             Ã   s.  �x| j dk�rt| j| jƒd }d}ytj| jjƒ |ƒI dH g}W n0 | jjtj	tj
fk
rt   | jjdƒ P Y nX |dgkr�| jjƒ  g }nZxXy|j| jjƒ ƒ 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�I dH }x|D ]}| jjƒ  �q4W |dk�r`| jjdƒ P |jdk �sx|jdk�r| jjd|jƒ | jƒ  P qy\xV|D ]N}|j�rÄ| jj|jdd�ƒI dH  n| jj|jdd�ƒI dH  | jjƒ  �qœW W q tj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.
        r(   rœ   Nzpacket queue is empty, abortingr   r   )r\   ÚPOSTzContent-Typezapplication/octet-stream)r�   r   rF   z*Connection refused by the server, abortingrG   i,  z6Unexpected status code %s in server response, abortingF)rp   z5Write loop: WebSocket connection was closed, abortingzExiting write loop taskéÿÿÿÿr¡   )"r   r�   r`   ra   r5   r›   r   Úgetr>   r�   r–   rN   rO   Ú	task_donerb   Ú
get_nowaitr.   r   rY   rS   rQ   rz   rU   rx   rW   r3   r#   r/   Z
send_bytesr{   rM   rw   rŸ   rR   )r	   rF   r\   rk   rj   rm   r   r   r
   rg     sj    







zAsyncClient._write_loop)N)F)r   )NNN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r!   r'   r)   r:   r;   r   rA   r3   rn   rd   rc   r$   rS   r-   re   ri   r€   rg   Ú__classcell__r   r   )rD   r
   r      s.   %


7b 
%+$r   )r5   ro   rM   ÚImportErrorr   Ú r   r   r   r   ZClientr   r   r   r   r
   Ú<module>   s   
