Ë
    Ùqj:2  ã                  óÂ   — d dl mZ d dlZd dlZd dlZd dlmZmZmZm	Z	m
Z
mZ ddlmZ ddlmZmZmZmZ ddlmZ dd	lmZ d
gZ ej.                  d«      Z G d„ d
«      Zy)é    )ÚannotationsN)ÚAnyÚCallableÚIterableÚIteratorÚLiteralÚoverloadé   )ÚConcurrencyError)Ú	OP_BINARYÚOP_CONTÚOP_TEXTÚFrame)ÚDataé   )ÚDeadlineÚ	Assemblerzutf-8c                  ó  — e Zd ZdZddd„ d„ f	 	 	 	 	 	 	 	 	 dd„Zddd„Zdd„Zedd„«       Zedd	„«       Zeddd
„«       Zed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   ó    úY/var/www/html/kelly-backend/venv/lib/python3.12/site-packages/websockets/sync/messages.pyÚ<lambda>zAssembler.<lambda>&   s   € ¨4r   c                  ó   — y r   r   r   r   r   r   zAssembler.<lambda>'   s   € ¨Dr   c                ó8  — 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)Ú	threadingÚLockÚmutexÚqueueÚSimpleQueueÚframesÚ
ValueErrorÚhighÚlowÚpauseÚresumeÚpausedÚget_in_progressÚclosed)Úselfr%   r&   r'   r(   s        r   Ú__init__zAssembler.__init__"   s³   € ô —^‘^Ó%ˆŒ
ô 8=×7HÑ7HÓ7JˆŒð Ð  Ø˜!‘)ˆCØˆ<˜C˜OØ˜‘7ˆDØÐ  Ø�QŠwÜ Ð!HÓIÐIØ�cŠzÜ Ð!LÓMÐMØ" CÐˆŒ	�4”8ØˆŒ
ØˆŒØˆŒð  %ˆÔð ˆ�r   c                óŽ  — | j                   r	 | j                  j                  d¬«      }nB	 |�"|dk  r| j                  j                  d¬«      }n| j                  j                  d|¬«      }|€t        d«      ‚|S # t        j                  $ r t        d«      d ‚w xY w# t        j                  $ r t        d|d›d	�«      d ‚w xY w)
NF©Úblockústream of frames endedr   T)r0   Útimeoutztimed out in z.1fÚs)r+   r#   Úgetr!   ÚEmptyÚEOFErrorÚTimeoutError)r,   r2   Úframes      r   Úget_next_framezAssembler.get_next_frameH   sÎ   € ð �;Š;ðCØŸ™Ÿ™¨e˜Ó4‘ðMð Ð&¨7°aª<Ø ŸK™KŸO™O°%˜OÓ8‘Eà ŸK™KŸO™O°$À˜OÓH�Eð ˆ=ÜÐ3Ó4Ð4Øˆøô —;‘;ò CÜÐ7Ó8¸dÐBðCûô —;‘;ò MÜ" ]°7¸3°-¸qÐ#AÓBÈÐLðMús   ŽA< ¬AB Á< BÂ%Cc                ób  — | j                   5  g }	 	 |j                  | j                  j                  d¬«      «       Œ,# t        j
                  $ r Y nw xY w|D ]  }| j                  j                  |«       Œ |D ]  }| j                  j                  |«       Œ 	 d d d «       y # 1 sw Y   y xY w)NFr/   )r    Úappendr#   r4   r!   r5   Úput)r,   r#   Úqueuedr8   s       r   Úreset_queuezAssembler.reset_queue^   sŽ   € ð �Z‹ZØˆFðØØ—M‘M $§+¡+§/¡/¸ /Ó">Ô?ð øä—;‘;ò Ùðúã�Ø—‘—‘ Õ&ð  ó  �Ø—‘—‘ Õ&ñ  ÷ �Z‰Zús'   �B%‘->¾AÁB%ÁAÁAB%Â%B.c                 ó   — y r   r   ©r,   r2   Údecodes      r   r4   zAssembler.gett   s   € ØHKr   c                 ó   — y r   r   r@   s      r   r4   zAssembler.getw   s   € ØKNr   c                ó   — y r   r   r@   s      r   r4   zAssembler.getz   s   € ØRUr   c                ó   — y r   r   r@   s      r   r4   zAssembler.get}   ó   € ØUXr   c                 ó   — y r   r   r@   s      r   r4   zAssembler.get€   rE   r   c                ó˜  — | j                   5  | j                  rt        d«      ‚d| _        ddd«       	 t        |«      }| j	                  |j                  d¬«      «      }| j                   5  | j                  «        ddd«       |j                  t        u s|j                  t        u sJ ‚|€|j                  t        u }|g}|j                  sy	 | j	                  |j                  d¬«      «      }| j                   5  | j                  «        ddd«       |j                  t        u sJ ‚|j                  |«       |j                  sŒyd| _        dj                  d„ |D «       «      }|r|j!                  «       S |S # 1 sw Y   �ŒQxY w# 1 sw Y   �ŒxY w# t        $ r | j                  |«       ‚ w xY w# 1 sw Y   Œ§xY w# d| _        w xY w)a?  
        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:
            timeout: If a timeout is provided and elapses before a complete
                message is received, :meth:`get` raises :exc:`TimeoutError`.
            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.
            TimeoutError: If a timeout is provided and elapses before a
                complete message is received.

        ú&get() or get_iter() is already runningTNF)Úraise_if_elapsedr   c              3  ó4   K  — | ]  }|j                   –— Œ y ­wr   )Údata)Ú.0r8   s     r   Ú	<genexpr>z Assembler.get.<locals>.<genexpr>Å   s   è ø€ Ð7± u˜Ÿ
�
±ùs   ‚)r    r*   r   r   r9   r2   Úmaybe_resumeÚopcoder   r   Úfinr7   r>   r   r;   ÚjoinrA   )r,   r2   rA   Údeadliner8   r#   rK   s          r   r4   zAssembler.getƒ   s–  € ð4 �Z‹ZØ×#Ò#Ü&Ð'OÓPÐPØ#'ˆDÔ ÷ ð	)Ü Ó(ˆHð ×'Ñ'¨×(8Ñ(8È%Ð(8Ó(PÓQˆEØ—“Ø×!Ñ!Ô#÷ à—<‘<¤7Ñ*¨e¯l©l¼iÑ.GÐGÐGØˆ~ØŸ™¬Ð0�Ø�WˆFð —i’iðØ ×/Ñ/Ø ×(Ñ(¸%Ð(Ó@ó�Eð —Z“ZØ×%Ñ%Ô'÷  à—|‘|¤wÑ.Ð.Ð.Ø—‘˜eÔ$ð —i“ið  $)ˆDÔ ð �x‰xÑ7±Ó7Ó7ˆÙØ—;‘;“=Ð àˆK÷Y ‰Zú÷ ‘ûô $ò ð ×$Ñ$ VÔ,Øð	ú÷
  �Zûð $)ˆDÕ ús_   �E;µ8G  Á-FÁ>AG  Ã!F Ã1G  Ã=F4Ä9G  Å;FÆFÆG  ÆF1Æ1G  Æ4F=Æ9G  Ç 	G	c                 ó   — y r   r   ©r,   rA   s     r   Úget_iterzAssembler.get_iterË   s   € Ø@Cr   c                 ó   — y r   r   rT   s     r   rU   zAssembler.get_iterÎ   s   € ØCFr   c                 ó   — y r   r   rT   s     r   rU   zAssembler.get_iterÑ   s   € ØFIr   c              #  óŠ  K  — | j                   5  | j                  rt        d«      ‚d| _        ddd«       | j                  «       }| j                   5  | j	                  «        ddd«       |j
                  t        u s|j
                  t        u sJ ‚|€|j
                  t        u }|r3t        «       }|j                  |j                  |j                  «      –— nt        |j                  «      –— |j                  s˜| j                  «       }| j                   5  | j	                  «        ddd«       |j
                  t        u sJ ‚|r)j                  |j                  |j                  «      –— nt        |j                  «      –— |j                  sŒ˜d| _        y# 1 sw Y   �ŒqxY w# 1 sw Y   �ŒIxY w# 1 sw Y   ŒŽxY w­w)a©  
        Stream the next message.

        Iterating the return value of :meth:`get_iter` 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.

        rH   TNF)r    r*   r   r9   rN   rO   r   r   ÚUTF8DecoderrA   rK   rP   Úbytesr   )r,   rA   r8   Údecoders       r   rU   zAssembler.get_iterÔ   sV  è ø€ ð2 �Z‹ZØ×#Ò#Ü&Ð'OÓPÐPØ#'ˆDÔ ÷ ð ×#Ñ#Ó%ˆØ�Z‹ZØ×ÑÔ÷ à�|‰|œwÑ&¨%¯,©,¼)Ñ*CÐCÐCØˆ>Ø—\‘\¤WÐ,ˆFÙÜ!“mˆGØ—.‘. §¡¨U¯Y©YÓ7Ó7ô ˜Ÿ
™
Ó#Ò#ð —)’)Ø×'Ñ'Ó)ˆEØ—“Ø×!Ñ!Ô#÷ à—<‘<¤7Ñ*Ð*Ð*ÙØ—n‘n U§Z¡Z°·±Ó;Ó;ô ˜EŸJ™JÓ'Ò'ð —)“)ð  %ˆÕ÷K ‰Zú÷ ‰Zú÷ �üsS   ‚G�F®$GÁF*Á#B6GÄF7Ä*A*GÆGÆF'Æ"GÆ*F4Æ/GÆ7G Æ<Gc                óÊ   — | j                   5  | j                  rt        d«      ‚| j                  j	                  |«       | j                  «        ddd«       y# 1 sw Y   yxY w)z
        Add ``frame`` to the next message.

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

        r1   N)r    r+   r6   r#   r<   Úmaybe_pause)r,   r8   s     r   r<   zAssembler.put  sD   € ð �Z‹ZØ�{Š{ÜÐ7Ó8Ð8à�K‰K�O‰O˜EÔ"Ø×ÑÔ÷ �Z‰Zús   �AAÁA"c                óî   — | j                   €y| j                  j                  «       s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%   r    Úlockedr#   Úqsizer)   r'   ©r,   s    r   r]   zAssembler.maybe_pause-  s`   € ð �9‰9ÐØà�z‰z× Ñ Ô"Ð"Ð"ð �;‰;×ÑÓ §¡Ò*°4·;²;ØˆDŒKØ�J‰J�Lð 4?Ð*r   c                óî   — | j                   €y| j                  j                  «       s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&   r    r_   r#   r`   r)   r(   ra   s    r   rN   zAssembler.maybe_resume:  s`   € ð �8‰8ÐØà�z‰z× Ñ Ô"Ð"Ð"ð �;‰;×ÑÓ $§(¡(Ò*¨t¯{ª{ØˆDŒKØ�K‰K�Mð 0;Ð*r   c                ó  — | j                   5  | j                  r
	 ddd«       yd| _        | j                  r| j                  j	                  d«       | j
                  rd| _        | j                  «        ddd«       y# 1 sw Y   yxY w)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`.

        NTF)r    r+   r*   r#   r<   r)   r(   ra   s    r   ÚclosezAssembler.closeG  s_   € ð �Z‹ZØ�{Š{Ø÷ ˆZð ˆDŒKà×#Ò#à—‘—‘ Ô%à�{Š{à#�”Ø—‘”÷ �Z‰Zús   �A>¤AA>Á>B)
r%   ú
int | Noner&   re   r'   úCallable[[], Any]r(   rf   ÚreturnÚNoner   )r2   úfloat | Nonerg   r   )r#   zIterable[Frame]rg   rh   )r2   ri   rA   úLiteral[True]rg   Ústr)r2   ri   rA   úLiteral[False]rg   rZ   )NN)r2   ri   rA   úbool | Nonerg   r   )rA   rj   rg   zIterator[str])rA   rl   rg   zIterator[bytes])rA   rm   rg   zIterator[Data])r8   r   rg   rh   )rg   rh   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r-   r9   r>   r	   r4   rU   r<   r]   rN   rd   r   r   r   r   r      sê   „ ñð   ØÙ#/Ù$0ð$àð$ð ð$ð !ð	$ð
 "ð$ð 
ó$ôLó,'ð, ÚKó ØKàÚNó ØNàÛUó ØUàÛXó ØXàÛXó ØXôFðP ÚCó ØCàÚFó ØFàÛIó ØIô>%ó@ó2óôr   )Ú
__future__r   Úcodecsr!   r   Útypingr   r   r   r   r   r	   Ú
exceptionsr   r#   r   r   r   r   r   Úutilsr   Ú__all__ÚgetincrementaldecoderrY   r   r   r   r   Ú<module>ry      sM   ðÝ "ã Û Û ß G× Gå )ß 7Ó 7Ý Ý ð ˆ-€à*ˆf×*Ñ*¨7Ó3€÷Iò Ir   