o
    .j                     @  sz   d dl mZ d dlZd dlmZ ddlmZmZ ejr+d dl	Z	ddl
mZ ddlmZ G d	d
 d
eZG dd deZdS )    )annotationsN)Future   )ServerSentEvent ServerSideEventExtensionFromHTTP)HTTPResponse   )ASGISSEExtensionc                   @  s`   e Zd ZdZd ddZed!d	d
Zddd"ddZd#ddZd$ddZ	d%ddZ
d&ddZdS )'WSGISSEExtensionzxSSE extension for WSGI applications.

    Reads from a WSGI response iterator, buffers text, and parses SSE events.
    	generator#typing.Generator[bytes, None, None]returnNonec                 C  s   || _ d| _d| _d | _d S )NF )
_generator_closed_buffer_last_event_id)selfr    r   Q/home/thesage/.local/lib/python3.10/site-packages/niquests/extensions/sgi/_sse.py__init__   s   
zWSGISSEExtension.__init__boolc                 C  s   | j S N)r   r   r   r   r   closed   s   zWSGISSEExtension.closedFrawr   ServerSentEvent | str | Nonec                C  s   | j rtd	 | jd}|dkr"| jd}|dkrd}nd}nd}|dkrK| jd| }| j|| d | _| |}|durJ|rH|d S |S q|  }|du rXd| _ dS |  j|7  _q)	zdRead and parse the next SSE event from the WSGI response.
        Returns None when the stream ends.zThe SSE extension is closedTz

z

      N)r   OSErrorr   find_parse_event_read_chunk)r   r   sep_idxsep_len	raw_eventeventchunkr   r   r   next_payload   s2   
zWSGISSEExtension.next_payload
str | Nonec                 C  s6   zt | j}|r|dW S W dS  ty   Y dS w )z4Read the next chunk from the WSGI response iterator.zutf-8N)nextr   decodeStopIteration)r   r*   r   r   r   r%   G   s   
zWSGISSEExtension._read_chunkr(   strServerSentEvent | Nonec              
   C  s   i }|  D ]E}|r|drq|d\}}}|dvrq|dr(|dd }|dkr1d|v r1q|dkrGzt|}W n ttfyF   Y qw |||< q|sPdS d|vr^| jdur^| j|d< td	i |}|jrl|j| _|S )
z3Parse a raw SSE event block into a ServerSentEvent.:>   iddatar)   retry r   Nr3    r5   r   )	
splitlines
startswith	partitionint
ValueError	TypeErrorr   r   r3   )r   r(   kwargslinekey_valuer)   r   r   r   r$   Q   s6   


zWSGISSEExtension._parse_eventc                 C  s.   | j rdS d| _ t| jdr| j  dS dS )z'Close the stream and release resources.NTclose)r   hasattrr   rC   r   r   r   r   rC   t   s   zWSGISSEExtension.closeresponser   c                 C     t r   NotImplementedErrorr   rE   r   r   r   start~      zWSGISSEExtension.startN)r   r   r   r   r   r   r   r   r   r   )r   r,   )r(   r0   r   r1   r   r   rE   r   r   r   )__name__
__module____qualname____doc__r   propertyr   r+   r%   r$   rC   rJ   r   r   r   r   r
      s    

(


#
r
   c                   @  sL   e Zd ZdZddd	ZedddZdddddZdddZdddZ	dS )ThreadASGISSEExtensionzSynchronous SSE extension wrapping an async ASGISSEExtension.

    Delegates all operations to the async extension on a background event loop,
    blocking the calling thread via concurrent.futures.Future.
    	async_extr	   loopasyncio.AbstractEventLoopr   r   c                 C  s   || _ || _d S r   )
_async_ext_loop)r   rV   rW   r   r   r   r      s   
zThreadASGISSEExtension.__init__r   c                 C  s   | j jS r   )rY   r   r   r   r   r   r      s   zThreadASGISSEExtension.closedFr   r   r   c                  s6   t  dfdd j fdd  S )	z9Block until the next SSE event arrives from the ASGI app.r   r   c               
     sX   zj jdI d H }  |  W d S  ty+ } z | W Y d }~d S d }~ww )Nr   )rY   r+   
set_result	Exceptionset_exception)resulte)futurer   r   r   r   _do   s   z0ThreadASGISSEExtension.next_payload.<locals>._doc                        j   S r   rZ   create_taskr   ra   r   r   r   <lambda>       z5ThreadASGISSEExtension.next_payload.<locals>.<lambda>NrN   r   rZ   call_soon_threadsafer^   )r   r   r   )ra   r`   r   r   r   r+      s   z#ThreadASGISSEExtension.next_payloadrE   r   c                 C  rF   r   rG   rI   r   r   r   rJ      rK   zThreadASGISSEExtension.startc                   s8   t  dfdd j fdd   dS )	z"Close the SSE stream 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   )rY   rC   r[   r\   r]   )r_   )r`   r   r   r   ra      s   z)ThreadASGISSEExtension.close.<locals>._doc                     rb   r   rc   r   re   r   r   rf      rg   z.ThreadASGISSEExtension.close.<locals>.<lambda>NrN   rh   r   r   )ra   r`   r   r   rC      s   zThreadASGISSEExtension.closeN)rV   r	   rW   rX   r   r   rL   rM   rO   rN   )
rP   rQ   rR   rS   r   rT   r   r+   rJ   rC   r   r   r   r   rU      s    

rU   )
__future__r   typingconcurrent.futuresr   *packages.urllib3.contrib.webextensions.sser   r   TYPE_CHECKINGasynciopackages.urllib3r   _async._sser	   r
   rU   r   r   r   r   <module>   s    s