Ë
    Õ±1h?1  ã                  óÂ   — 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   ó    úV/var/www/skyplay_api_hub/venv/lib/python3.12/site-packages/websockets/sync/messages.pyú<lambda>zAssembler.<lambda>&   ó   � r   c                  ó   — y r   r   r   r   r   r   zAssembler.<lambda>'   r   r   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)r1   Útimeoutztimed out in z.1fÚs)r,   r$   Úgetr"   ÚEmptyÚEOFErrorÚTimeoutError)r-   r3   Ú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)NFr0   )r!   Úappendr$   r5   r"   r6   Úput)r-   r$   Úqueuedr9   s       r   Úreset_queuezAssembler.reset_queue^   s£   € ð �Z‰Zñ 	'ØˆFðØØ—M‘M $§+¡+§/¡/¸ /Ó">Ô?ð øä—;‘;ò Ùðúàò '�Ø—‘—‘ Õ&ð'ð  ò '�Ø—‘—‘ Õ&ñ'÷	'÷ 	'ñ 	'ús'   �B%‘->¾AÁB%ÁAÁAB%Â%B.c                 ó   — y r   r   ©r-   r3   Údecodes      r   r5   zAssembler.gett   s   € ØHKr   c                 ó   — y r   r   rA   s      r   r5   zAssembler.getw   s   € ØKNr   c                ó   — y r   r   rA   s      r   r5   zAssembler.getz   s   € ØRUr   c                ó   — y r   r   rA   s      r   r5   zAssembler.get}   ó   € ØUXr   c                 ó   — y r   r   rA   s      r   r5   zAssembler.get€   rF   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)Ú.0r9   s     r   ú	<genexpr>z Assembler.get.<locals>.<genexpr>Ä   s   è ø€ Ò7 u˜Ÿ
�
Ñ7ùs   ‚)r!   r+   r   r   r:   r3   Úmaybe_resumeÚopcoder   r   Úfinr8   r?   r   r<   ÚjoinrB   )r-   r3   rB   Údeadliner9   r$   rL   s          r   r5   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÷W	(ñ 	(ú÷$ñ $ûô $ò ð ×$Ñ$ VÔ,Øð	ú÷
(ð (ûð $)ˆ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-   rB   s     r   Úget_iterzAssembler.get_iterÊ   s   € Ø@Cr   c                 ó   — y r   r   rU   s     r   rV   zAssembler.get_iterÍ   s   € ØCFr   c                 ó   — y r   r   rU   s     r   rV   zAssembler.get_iterÐ   s   € ØFIr   c              #  óf  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|j                  –— |j                  s�| j                  «       }| j                   5  | j	                  «        ddd«       |j
                  t        u sJ ‚|r)j                  |j                  |j                  «      –— n|j                  –— |j                  sŒ�d| _        y# 1 sw Y   �Œ_xY w# 1 sw Y   �Œ7x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.

        rI   TNF)r!   r+   r   r:   rO   rP   r   r   ÚUTF8DecoderrB   rL   rQ   r   )r-   rB   r9   Údecoders       r   rV   zAssembler.get_iterÓ   s`  è ø€ ð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°·±Ó;Ó;à—j‘jÒ ð —)“)ð  %ˆÕ÷G	(ñ 	(ú÷	 ñ 	 ú÷$ð $üsS   ‚F1�F®$F1ÁFÁ#B-F1ÄF%Ä!A!F1ÆF1ÆFÆF1ÆF"ÆF1Æ%F.Æ*F1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.

        r2   N)r!   r,   r7   r$   r=   Úmaybe_pause)r-   r9   s     r   r=   zAssembler.put  sO   € ð �Z‰Zñ 	Ø�{Š{ÜÐ7Ó8Ð8à�K‰K�O‰O˜EÔ"Ø×ÑÔ÷	÷ 	ñ 	ú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   rO   zAssembler.maybe_resume7  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.closeD  sm   € ð �Z‰Zñ 	Ø�{Š{Ø÷	ð 	ð ˆDŒKà×#Ò#à—‘—‘ Ô%à�{Š{à#�”Ø—‘”÷	÷ 	ñ 	ús   �A>¤AA>Á>B)
r&   ú
int | Noner'   re   r(   úCallable[[], Any]r)   rf   ÚreturnÚNoner   )r3   úfloat | Nonerg   r   )r$   zIterable[Frame]rg   rh   )r3   ri   rB   úLiteral[True]rg   Ústr)r3   ri   rB   úLiteral[False]rg   Úbytes)NN)r3   ri   rB   úbool | Nonerg   r   )rB   rj   rg   zIterator[str])rB   rl   rg   zIterator[bytes])rB   rn   rg   zIterator[Data])r9   r   rg   rh   )rg   rh   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r.   r:   r?   r	   r5   rV   r=   r]   rO   rd   r   r   r   r   r      sê   „ ñð   ØÙ#/Ù$0ð$àð$ð ð$ð !ð	$ð
 "ð$ð 
ó$ôLó,'ð, ÚKó ØKàÚNó ØNàÛUó ØUàÛXó ØXàÛXó ØXôEðN Ú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__ÚgetincrementaldecoderrZ   r   r   r   r   ú<module>rz      sM   ðÝ "ã Û Û ß G× Gå )ß 7Ó 7Ý Ý ð ˆ-€à*ˆf×*Ñ*¨7Ó3€÷Fò Fr   