o
    .j	                     @  s   d dl mZ d dlZd dlmZ zddlmZ W n ey'   ddlm	Z Y nw ej
r;d dlZddlmZ dd	lmZ G d
d deZdS )    )annotationsN)Future   )WebSocketExtensionFromHTTP)RawExtensionFromHTTP)HTTPResponse   )ASGIWebSocketExtensionc                   @  sP   e Zd ZdZddd	ZdddZedddZdddZdddZ	d ddZ
dS )!ThreadASGIWebSocketExtensionzSynchronous WebSocket extension wrapping an async ASGIWebSocketExtension.

    Delegates all operations to the async extension on a background event loop,
    blocking the calling thread via concurrent.futures.Future.
    	async_extr	   loopasyncio.AbstractEventLoopreturnNonec                 C  s   || _ || _d S N)
_async_ext_loop)selfr   r    r   P/home/thesage/.local/lib/python3.10/site-packages/niquests/extensions/sgi/_ws.py__init__   s   
z%ThreadASGIWebSocketExtension.__init__responser   c                 C  s   t r   )NotImplementedError)r   r   r   r   r   start   s   z"ThreadASGIWebSocketExtension.startboolc                 C  s   | j jS r   )r   closedr   r   r   r   r       s   z#ThreadASGIWebSocketExtension.closedstr | bytes | Nonec                   s4   t  dfdd j fdd  S )	zABlock until the next message arrives from the ASGI WebSocket app.r   r   c               
     sT   zj  I d H }  |  W d S  ty) } z | W Y d }~d S d }~ww r   )r   next_payload
set_result	Exceptionset_exception)resultefuturer   r   r   _do(      z6ThreadASGIWebSocketExtension.next_payload.<locals>._doc                        j   S r   r   create_taskr   r&   r   r   r   <lambda>/       z;ThreadASGIWebSocketExtension.next_payload.<locals>.<lambda>Nr   r   r   r   call_soon_threadsafer"   r   r   r&   r%   r   r   r   $   s   z)ThreadASGIWebSocketExtension.next_payloadbufstr | bytesc                   s:   t  dfdd j fdd   dS )	z)Send a message to the ASGI WebSocket app.r   r   c               
     sV   zj  I d H  d  W d S  ty* }  z|  W Y d } ~ d S d } ~ ww r   )r   send_payloadr   r    r!   r#   )r2   r%   r   r   r   r&   6   s   z6ThreadASGIWebSocketExtension.send_payload.<locals>._doc                     r(   r   r)   r   r+   r   r   r,   =   r-   z;ThreadASGIWebSocketExtension.send_payload.<locals>.<lambda>Nr.   r/   )r   r2   r   )r&   r2   r%   r   r   r4   2   s   z)ThreadASGIWebSocketExtension.send_payloadc                   s8   t  dfdd j fdd   dS )	z!Close the WebSocket and clean up.r   r   c               
     sT   zj  I d H   d  W d S  ty) }  z |  W Y d } ~ d S d } ~ ww r   )r   closer   r    r!   r5   r$   r   r   r&   D   r'   z/ThreadASGIWebSocketExtension.close.<locals>._doc                     r(   r   r)   r   r+   r   r   r,   K   r-   z4ThreadASGIWebSocketExtension.close.<locals>.<lambda>Nr.   r/   r   r   r1   r   r6   @   s   z"ThreadASGIWebSocketExtension.closeN)r   r	   r   r   r   r   )r   r   r   r   )r   r   )r   r   )r2   r3   r   r   r.   )__name__
__module____qualname____doc__r   r   propertyr   r   r4   r6   r   r   r   r   r
      s    



r
   )
__future__r   typingconcurrent.futuresr   )packages.urllib3.contrib.webextensions.wsr   ImportError*packages.urllib3.contrib.webextensions.rawr   TYPE_CHECKINGasynciopackages.urllib3r   
_async._wsr	   r
   r   r   r   r   <module>   s    