Ë
    ¸�©jì(  ã                  óº   — d dl mZ d dlZd dlZd dlmZ d dlmZmZm	Z	m
Z
 d dlZddlmZ ddlmZmZmZmZ ddlmZ d	gZ ej*                  d
«      Z G d„ d	«      Zy)é    )ÚannotationsN)ÚAsyncIterator)ÚAnyÚCallableÚLiteralÚoverloadé   )ÚConcurrencyError)ÚBINARYÚCONTÚTEXTÚFrame)ÚDataÚ	Assemblerzutf-8c                  óÜ   — e Zd ZdZddd„ d„ f	 	 	 	 	 	 	 	 	 dd„Zedd„«       Zedd„«       Zeddd„«       Zddd	„Zedd
„«       Zedd„«       Zeddd„«       Zddd„Zdd„Zdd„Z	dd„Z
dd„Zy)r   aË  
    Assemble messages from frames.

    :class:`Assembler` expects only data frames. The stream of frames must
    respect the protocol; if it doesn't, the behavior is undefined.

    Args:
        pause: Called when the buffer of frames goes above the high water mark;
            should pause reading from the network.
        resume: Called when the buffer of frames goes below the low water mark;
            should resume reading from the network.

    Nc                  ó   — y ©N© r   ó    ú_/var/www/html/streamforge-license/venv/lib/python3.12/site-packages/websockets/trio/messages.pyú<lambda>zAssembler.<lambda>'   ó   � r   c                  ó   — y r   r   r   r   r   r   zAssembler.<lambda>(   r   r   c                ó<  — |  |  t        j                  t        j                  «      \  | _        | _        |�|€|dz  }|€|�|dz  }|�"|� |dk  rt        d«      ‚||k  rt        d«      ‚||c| _        | _        || _	        || _
        d| _        d| _        d| _        y )Né   r   z%low must be positive or equal to zeroz)high must be greater than or equal to lowF)ÚtrioÚopen_memory_channelÚmathÚinfÚsend_framesÚrecv_framesÚ
ValueErrorÚhighÚlowÚpauseÚresumeÚpausedÚget_in_progressÚclosed)Úselfr#   r$   r%   r&   s        r   Ú__init__zAssembler.__init__#   s¸   € ñ 	ÙÜ-1×-EÑ-EÄdÇhÁhÓ-OÑ*ˆÔ˜$Ô*ð Ð  Ø˜!‘)ˆCØˆ<˜C˜OØ˜‘7ˆDØÐ  Ø�QŠwÜ Ð!HÓIÐIØ�cŠzÜ Ð!LÓMÐMØ" CÐˆŒ	�4”8ØˆŒ
ØˆŒØˆŒð  %ˆÔð ˆ�r   c              ƒ  ó   K  — y ­wr   r   ©r*   Údecodes     r   ÚgetzAssembler.getG   s	   è ø€ Ø7:ùó   ‚c              ƒ  ó   K  — y ­wr   r   r-   s     r   r/   zAssembler.getJ   s	   è ø€ Ø:=ùr0   c              ƒ  ó   K  — y ­wr   r   r-   s     r   r/   zAssembler.getM   s	   è ø€ Ø=@ùr0   c              ƒ  óú  K  — | j                   rt        d«      ‚d| _         	 	 | j                  j                  «       ƒ d{  –—† }| j                  «        |j                  t        u s|j                  t        u sJ ‚|€|j                  t        u }|g}|j                  se	 | j                  j                  «       ƒ d{  –—† }| j                  «        |j                  t$        u sJ ‚|j'                  |«       |j                  sŒed| _         dj)                  d	„ |D «       «      }|r|j+                  «       S |S 7 Œõ# t        j
                  $ r t        d«      ‚w xY w7 Œ�# t        j                  $ r` | j                  j                  }|j                  rJ d«       ‚|j                   rJ d«       ‚|D ]  }| j                  j#                  |«       Œ ‚ t        j
                  $ r t        d«      ‚w xY w# d| _         w xY w­w)
a0  
        Read the next message.

        :meth:`get` returns a single :class:`str` or :class:`bytes`.

        If the message is fragmented, :meth:`get` waits until the last frame is
        received, then it reassembles the message and returns it. To receive
        messages frame by frame, use :meth:`get_iter` instead.

        Args:
            decode: :obj:`False` disables UTF-8 decoding of text frames and
                returns :class:`bytes`. :obj:`True` forces UTF-8 decoding of
                binary frames and returns :class:`str`.

        Raises:
            EOFError: If the stream of frames has ended.
            UnicodeDecodeError: If a text frame contains invalid UTF-8.
            ConcurrencyError: If two coroutines run :meth:`get` or
                :meth:`get_iter` concurrently.

        ú&get() or get_iter() is already runningTNústream of frames endedzno task should receivezqueue should be emptyFr   c              3  ó4   K  — | ]  }|j                   –— Œ y ­wr   )Údata)Ú.0Úframes     r   ú	<genexpr>z Assembler.get.<locals>.<genexpr>‘   s   è ø€ Ò7 u˜Ÿ
�
Ñ7ùs   ‚)r(   r
   r!   Úreceiver   ÚEndOfChannelÚEOFErrorÚmaybe_resumeÚopcoder   r   ÚfinÚ	Cancelledr    Ú_stateÚreceive_tasksr7   Úsend_nowaitr   ÚappendÚjoinr.   )r*   r.   r9   ÚframesÚstater7   s         r   r/   zAssembler.getP   sÜ  è ø€ ð, ×ÒÜ"Ð#KÓLÐLØ#ˆÔð
!	)ð9Ø"×.Ñ.×6Ñ6Ó8×8�ð ×ÑÔØ—<‘<¤4Ñ'¨5¯<©<¼6Ñ+AÐAÐAØˆ~ØŸ™¬Ð-�Ø�WˆFð —i’ið=Ø"&×"2Ñ"2×":Ñ":Ó"<×<�Eð ×!Ñ!Ô#Ø—|‘|¤tÑ+Ð+Ð+Ø—‘˜eÔ$ð# —i“ið( $)ˆDÔ ð �x‰xÑ7°Ô7Ó7ˆÙØ—;‘;“=Ð àˆKðK 9ùÜ×$Ñ$ò 9ÜÐ7Ó8Ð8ð9úð =ùÜ—~‘~ò 	ð !×,Ñ,×3Ñ3�EØ$×2Ò2ÐLÐ4LÓLÐ2Ø$ŸzšzÐBÐ+BÓB˜>Ø!'ò <˜Ø×(Ñ(×4Ñ4°UÕ;ð<àÜ×(Ñ(ò =Ü"Ð#;Ó<Ð<ð=ûð $)ˆDÕ üsm   ‚G;£D8 Á D6ÁD8 ÁAG/ ÂE Â<EÂ=E ÃAG/ Ä3G;Ä6D8 Ä8EÅG/ ÅE ÅBG,Ç,G/ Ç/	G8Ç8G;c                 ó   — y r   r   r-   s     r   Úget_iterzAssembler.get_iter—   s   € ØEHr   c                 ó   — y r   r   r-   s     r   rJ   zAssembler.get_iterš   s   € ØHKr   c                 ó   — y r   r   r-   s     r   rJ   zAssembler.get_iter�   s   € ØKNr   c               óÚ  K  — | j                   rt        d«      ‚d| _         	 | j                  j                  «       ƒ d{  –—† }| j                  «        |j                  t        u s|j                  t        u sJ ‚|€|j                  t        u }|r4t        «       }|j                  |j                  |j                  «      ­–— nt!        |j                  «      ­–— |j                  s˜	 | j                  j                  «       ƒ d{  –—† }| j                  «        |j                  t"        u sJ ‚|r*j                  |j                  |j                  «      ­–— nt!        |j                  «      ­–— |j                  sŒ˜d| _         y7 �ŒI# t        j
                  $ r	 d| _         ‚ t        j                  $ r t        d«      ‚w xY w7 ŒÀ# t        j                  $ r t        d«      ‚w xY w­w)a¸  
        Stream the next message.

        Iterating the return value of :meth:`get_iter` asynchronously yields a
        :class:`str` or :class:`bytes` for each frame in the message.

        The iterator must be fully consumed before calling :meth:`get_iter` or
        :meth:`get` again. Else, :exc:`ConcurrencyError` is raised.

        This method only makes sense for fragmented messages. If messages aren't
        fragmented, use :meth:`get` instead.

        Args:
            decode: :obj:`False` disables UTF-8 decoding of text frames and
                returns :class:`bytes`. :obj:`True` forces UTF-8 decoding of
                binary frames and returns :class:`str`.

        Raises:
            EOFError: If the stream of frames has ended.
            UnicodeDecodeError: If a text frame contains invalid UTF-8.
            ConcurrencyError: If two coroutines run :meth:`get` or
                :meth:`get_iter` concurrently.

        r4   TNFr5   )r(   r
   r!   r;   r   rA   r<   r=   r>   r?   r   r   ÚUTF8Decoderr.   r7   r@   Úbytesr   )r*   r.   r9   Údecoders       r   rJ   zAssembler.get_iter    s›  è ø€ ð2 ×ÒÜ"Ð#KÓLÐLØ#ˆÔð	5Ø×*Ñ*×2Ñ2Ó4×4ˆEð 	×ÑÔØ�|‰|œtÑ# u§|¡|´vÑ'=Ð=Ð=Øˆ>Ø—\‘\¤TÐ)ˆFÙÜ!“mˆGØ—.‘. §¡¨U¯Y©YÓ7Ô7ô ˜Ÿ
™
Ó#Ó#ð —)’)ð
9Ø"×.Ñ.×6Ñ6Ó8×8�ð ×ÑÔØ—<‘<¤4Ñ'Ð'Ð'ÙØ—n‘n U§Z¡Z°·±Ó;Ô;ô ˜EŸJ™JÓ'Ó'ð —)“)ð"  %ˆÕðG 5úÜ�~‰~ò 	Ø#(ˆDÔ ØÜ× Ñ ò 	5ÜÐ3Ó4Ð4ð	5úð( 9ùÜ×$Ñ$ò 9ÜÐ7Ó8Ð8ð9üs_   ‚G+¢F ¿FÁ F ÁB$G+Ã)G	 ÄGÄG	 ÄA4G+Æ G+ÆF Æ9GÇG+ÇG	 Ç	G(Ç(G+c                óˆ   — | j                   rt        d«      ‚| j                  j                  |«       | j	                  «        y)z
        Add ``frame`` to the next message.

        Raises:
            EOFError: If the stream of frames has ended.

        r5   N)r)   r=   r    rD   Úmaybe_pause)r*   r9   s     r   ÚputzAssembler.putê   s7   € ð �;Š;ÜÐ3Ó4Ð4à×Ñ×$Ñ$ UÔ+Ø×ÑÕr   c                óÔ   — | j                   €yt        | j                  j                  j                  «      | j                   kD  r%| j
                  sd| _        | j                  «        yyy)z7Pause the writer if queue is above the high water mark.NT)r#   Úlenr    rB   r7   r'   r%   ©r*   s    r   rR   zAssembler.maybe_pauseø   sV   € ð �9‰9ÐØô ˆt×Ñ×&Ñ&×+Ñ+Ó,¨t¯y©yÒ8ÀÇÂØˆDŒKØ�J‰J�Lð BMÐ8r   c                óÔ   — | j                   €yt        | j                  j                  j                  «      | j                   k  r%| j
                  rd| _        | j                  «        yyy)z7Resume the writer if queue is below the low water mark.NF)r$   rU   r    rB   r7   r'   r&   rV   s    r   r>   zAssembler.maybe_resume  sU   € ð �8‰8ÐØô ˆt×Ñ×&Ñ&×+Ñ+Ó,°·±Ò8¸T¿[º[ØˆDŒKØ�K‰K�Mð >IÐ8r   c                ó€   — | j                   ryd| _         | j                  j                  «        d„ | _        d„ | _        y)z½
        End the stream of frames.

        Calling :meth:`close` concurrently with :meth:`get`, :meth:`get_iter`,
        or :meth:`put` is safe. They will raise :exc:`EOFError`.

        NTc                  ó   — y r   r   r   r   r   r   z!Assembler.close.<locals>.<lambda>"  r   r   c                  ó   — y r   r   r   r   r   r   z!Assembler.close.<locals>.<lambda>#  r   r   )r)   r    Úcloser%   r&   rV   s    r   r[   zAssembler.close  s9   € ð �;Š;ØàˆŒð 	×Ñ×ÑÔ ñ "ˆŒ
Ù"ˆ�r   )
r#   ú
int | Noner$   r\   r%   úCallable[[], Any]r&   r]   ÚreturnÚNone)r.   úLiteral[True]r^   Ústr)r.   úLiteral[False]r^   rO   r   )r.   úbool | Noner^   r   )r.   r`   r^   zAsyncIterator[str])r.   rb   r^   zAsyncIterator[bytes])r.   rc   r^   zAsyncIterator[Data])r9   r   r^   r_   )r^   r_   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r+   r   r/   rJ   rS   rR   r>   r[   r   r   r   r   r      sÅ   „ ñð   ØÙ#/Ù$0ð"àð"ð ð"ð !ð	"ð
 "ð"ð 
ó"ðH Ú:ó Ø:àÚ=ó Ø=àÛ@ó Ø@ôEðN ÚHó ØHàÚKó ØKàÛNó ØNôH%óTó
ó
ô#r   )Ú
__future__r   Úcodecsr   Úcollections.abcr   Útypingr   r   r   r   r   Ú
exceptionsr
   rG   r   r   r   r   r   Ú__all__ÚgetincrementaldecoderrN   r   r   r   r   ú<module>ro      sM   ðÝ "ã Û Ý )ß 3Ó 3ã å )ß .Ó .Ý ð ˆ-€à*ˆf×*Ñ*¨7Ó3€÷O#ò O#r   