o
    .j                     @  sP   d dl mZ d dlZd dlZd dlZddlmZ ddlmZ G dd deZ	dS )    )annotationsN   )%AsyncServerSideEventExtensionFromHTTP)ServerSentEventc                   @  sd   e Zd ZdZd#ddZ	d$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"S )*ASGISSEExtensionzAsync SSE extension for ASGI applications.

    Runs the ASGI app as a normal HTTP streaming request and parses
    SSE events from the response body chunks.
    returnNonec                 C  s"   d| _ d| _d | _d | _d | _d S )NF )_closed_buffer_last_event_id_response_queue_taskself r   X/home/thesage/.local/lib/python3.10/site-packages/niquests/extensions/sgi/_async/_sse.py__init__   s
   
zASGISSEExtension.__init__    app
typing.Anyscopedict[str, typing.Any]bodybytesc                   s   t  _dt  dfdddfdd	d fd
d}t | _	 j I dH }|du r@td|d dkrH|S q0)zStart the ASGI app and wait for http.response.start.

        Returns the response start message (with status and headers).
        Fr   r   c                     s,   r  I d H  ddiS dd ddS )Ntypezhttp.disconnectTzhttp.requestF)r   r   	more_body)waitr   )r   request_completeresponse_completer   r   receive(   s   z'ASGISSEExtension.start.<locals>.receivemessager   c                   s@   j | I d H  | d dkr| dds   d S d S d S )Nr   http.response.bodyr   F)r   putgetset)r!   )r   r   r   r   send0   s
   z$ASGISSEExtension.start.<locals>.sendc                	     sB   z I d H  W j d I d H  d S j d I d H  w N)r   r#   r   )r   r    r   r   r&   r   r   run_app5   s   *z'ASGISSEExtension.start.<locals>.run_appTNz/ASGI app closed before sending response headersr   zhttp.response.start)r   r   )r!   r   r   r   r   r   )asyncioQueuer   Eventcreate_taskr   r$   ConnectionError)r   r   r   r   r(   r!   r   )r   r   r    r   r   r   r   r&   r   start   s   

zASGISSEExtension.startboolc                 C  s   | j S r'   )r
   r   r   r   r   closedE   s   zASGISSEExtension.closedF)rawr2   ServerSentEvent | str | Nonec                  s   | j rtd	 | jd}|dkr#| jd}|dkr d}nd}nd}|dkrL| jd| }| j|| d | _| |}|durK|rI|d S |S q|  I dH }|du r\d| _ dS |  j|7  _q	)	zkRead and parse the next SSE event from the ASGI response stream.
        Returns None when the stream ends.zThe SSE extension is closedTz

z

r      N)r
   OSErrorr   find_parse_event_read_chunk)r   r2   sep_idxsep_len	raw_eventeventchunkr   r   r   next_payloadI   s4   
zASGISSEExtension.next_payload
str | Nonec                   sV   | j du rdS | j  I dH }|du rdS |d dkr)|dd}|r)|dS dS )z0Read the next body chunk from the ASGI response.Nr   r"   r   r   zutf-8)r   r$   decode)r   r!   r   r   r   r   r9   q   s   

zASGISSEExtension._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    NrE    rG   r   )	
splitlines
startswith	partitionint
ValueError	TypeErrorr   r   rE   )r   r<   kwargslinekey_valuer=   r   r   r   r8      s6   


zASGISSEExtension._parse_eventc                   sv   | j rdS d| _ | jdur7| j s9| j  ttj | jI dH  W d   dS 1 s0w   Y  dS dS dS )z/Close the SSE stream and clean up the app task.NT)r
   r   donecancel
contextlibsuppressr*   CancelledErrorr   r   r   r   close   s   
"zASGISSEExtension.closeNr)   )r   )r   r   r   r   r   r   r   r   )r   r0   )r2   r0   r   r3   )r   r@   )r<   rB   r   rC   )__name__
__module____qualname____doc__r   r/   propertyr1   r?   r9   r8   r[   r   r   r   r   r      s    
,
(
#r   )

__future__r   r*   rX   typing1packages.urllib3.contrib.webextensions._async.sser   *packages.urllib3.contrib.webextensions.sser   r   r   r   r   r   <module>   s    