o
    .j~                     @  sL   d dl mZ d dlZd dlmZ ddlmZ ddlmZ G dd deZ	dS )	    )annotationsN)pyfetch   )%AsyncServerSideEventExtensionFromHTTP)ServerSentEventc                   @  sb   e Zd ZdZdddZd d!ddZed"ddZd#ddZddd$ddZ	d%ddZ
dddZdS )&AsyncPyodideSSEExtensionzAsync SSE extension for Pyodide using pyfetch streaming + manual SSE parsing.

    Reads from a ReadableStream reader, buffers partial lines,
    and parses complete SSE events.
    returnNonec                 C  s   d| _ d| _d | _d | _d S )NF )_closed_buffer_last_event_id_readerself r   \/home/thesage/.local/lib/python3.10/site-packages/niquests/extensions/pyodide/_async/_sse.py__init__   s   
z!AsyncPyodideSSEExtension.__init__Nurlstrheadersdict[str, str] | Nonec                   sR   dddd|p	i d}t |fi |I dH }|jj}|dur'| | _dS dS )z Open the SSE stream via pyfetch.GETztext/event-streamzno-store)AcceptzCache-Control)methodr   N)r   js_responsebody	getReaderr   )r   r   r   fetch_optionsr   r   r   r   r   start   s   	zAsyncPyodideSSEExtension.startboolc                 C  s   | j S N)r   r   r   r   r   closed)   s   zAsyncPyodideSSEExtension.closed
str | Nonec                   sj   | j du rdS z"| j  I dH }|jrW dS |j}|dur(t| dW S W dS  ty4   Y dS w )z,Read the next chunk from the ReadableStream.Nzutf-8)r   readdonevaluebytesto_pydecode	Exception)r   resultr&   r   r   r   _read_chunk-   s   
z$AsyncPyodideSSEExtension._read_chunkF)rawr-   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	)	z]Read and parse the next SSE event from the stream.
        Returns None when the stream ends.zThe SSE extension is closedTz

z

r      N)r   OSErrorr   find_parse_eventr,   )r   r-   sep_idxsep_len	raw_eventeventchunkr   r   r   next_payload=   s4   
z%AsyncPyodideSSEExtension.next_payloadr6   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datar7   retry    Nr<    r>   r   )	
splitlines
startswith	partitionint
ValueError	TypeErrorr   r   r<   )r   r6   kwargslinekey_r&   r7   r   r   r   r3   f   s6   


z%AsyncPyodideSSEExtension._parse_eventc                   sR   | j rdS d| _ | jdur'z
| j I dH  W n	 ty!   Y nw d| _dS dS )z'Close the stream and release resources.NT)r   r   cancelr*   r   r   r   r   close   s   

zAsyncPyodideSSEExtension.close)r   r	   r!   )r   r   r   r   r   r	   )r   r    )r   r#   )r-   r    r   r.   )r6   r   r   r:   )__name__
__module____qualname____doc__r   r   propertyr"   r,   r9   r3   rM   r   r   r   r   r      s    


)#r   )

__future__r   typingpyodide.httpr   1packages.urllib3.contrib.webextensions._async.sser   *packages.urllib3.contrib.webextensions.sser   r   r   r   r   r   <module>   s    