o
    .j                     @  sb   d dl mZ d dlZd dlmZ d dlmZ ddlmZm	Z	 ej
r'ddlmZ G dd	 d	e	ZdS )
    )annotationsN)run_sync)pyfetch   )ServerSentEvent ServerSideEventExtensionFromHTTP)HTTPResponsec                   @  sb   e Zd Z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(dd ZdS ))PyodideSSEExtensionzSSE extension for Pyodide using pyfetch streaming + manual SSE parsing.

    Synchronous via JSPI (run_sync). Reads from a ReadableStream reader,
    buffers partial lines, and parses complete SSE events.
    Nurlstrheadersdict[str, str] | NonereturnNonec                 C  sf   d| _ d| _d | _d | _dddd|pi d}tt|fi |}|jj}|d ur1| | _d S d S )NF GETztext/event-streamzno-store)AcceptzCache-Control)methodr   )	_closed_buffer_last_event_id_readerr   r   js_responsebody	getReader)selfr
   r   fetch_optionsr   r    r   U/home/thesage/.local/lib/python3.10/site-packages/niquests/extensions/pyodide/_sse.py__init__   s    	zPyodideSSEExtension.__init__boolc                 C  s   | j S N)r   r   r   r   r   closed*   s   zPyodideSSEExtension.closed
str | Nonec                 C  sf   | j du rdS z!t| j  }|jrW dS |j}|dur&t| dW S W dS  ty2   Y dS w )z?Read the next chunk from the ReadableStream, blocking via JSPI.Nzutf-8)	r   r   readdonevaluebytesto_pydecode	Exception)r   resultr'   r   r   r   _read_chunk.   s   
zPyodideSSEExtension._read_chunkF)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)	z]Read and parse the next SSE event from the stream.
        Returns None when the stream ends.zThe SSE extension is closedTz

z

      N)r   OSErrorr   find_parse_eventr-   )r   r.   sep_idxsep_len	raw_eventeventchunkr   r   r   next_payload>   s2   
z PyodideSSEExtension.next_payloadr8   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datar9   retry    Nr>    r@   r   )	
splitlines
startswith	partitionint
ValueError	TypeErrorr   r   r>   )r   r8   kwargslinekey_r'   r9   r   r   r   r5   g   s6   


z PyodideSSEExtension._parse_eventresponser   c                 C  s   t r!   )NotImplementedError)r   rN   r   r   r   start   s   zPyodideSSEExtension.startc                 C  sN   | j rdS d| _ | jdur%z	t| j  W n	 ty   Y nw d| _dS dS )z'Close the stream and release resources.NT)r   r   r   cancelr+   r"   r   r   r   close   s   

zPyodideSSEExtension.closer!   )r
   r   r   r   r   r   )r   r    )r   r$   )r.   r    r   r/   )r8   r   r   r<   )rN   r   r   r   )r   r   )__name__
__module____qualname____doc__r   propertyr#   r-   r;   r5   rP   rR   r   r   r   r   r	      s    

)
#r	   )
__future__r   typingpyodide.ffir   pyodide.httpr   *packages.urllib3.contrib.webextensions.sser   r   TYPE_CHECKINGpackages.urllib3r   r	   r   r   r   r   <module>   s    