ó
    ÖÁ“iy+  ã                  óò   • S SK Jr  S SKrS SKrS SKrS SKJrJr  S SKJ	r	J
r
JrJrJrJr  SSKJr  SSKJrJrJrJr  SSKJr  S	/r\R0                  " S
5      r\" S5      r " S S\\   5      r " S S	5      rg)é    )ÚannotationsN)ÚAsyncIteratorÚIterable)ÚAnyÚCallableÚGenericÚLiteralÚTypeVarÚoverloadé   )ÚConcurrencyError)Ú	OP_BINARYÚOP_CONTÚOP_TEXTÚFrame)ÚDataÚ	Assemblerzutf-8ÚTc                  óX   • \ rS rSrSrSS jrSS jrSS jrSSS jjrSS jr	SS jr
S	rg
)ÚSimpleQueueé   zy
Simplified version of :class:`asyncio.Queue`.

Provides only the subset of functionality needed by :class:`Assembler`.

c                óz   • [         R                  " 5       U l        S U l        [        R
                  " 5       U l        g ©N)ÚasyncioÚget_running_loopÚloopÚ
get_waiterÚcollectionsÚdequeÚqueue©Úselfs    ÚX/home/mande/repo/quber/.venv/lib/python3.13/site-packages/websockets/asyncio/messages.pyÚ__init__ÚSimpleQueue.__init__   s)   € Ü×,Ò,Ó.ˆŒ	Ø7;ˆŒÜ+6×+<Ò+<Ó+>ˆ�
ó    c                ó,   • [        U R                  5      $ r   )Úlenr    r!   s    r#   Ú__len__ÚSimpleQueue.__len__"   s   € Ü�4—:‘:‹Ðr&   c                óÌ   • U R                   R                  U5        U R                  b<  U R                  R                  5       (       d  U R                  R	                  S5        ggg)zPut an item into the queue.N)r    Úappendr   ÚdoneÚ
set_result)r"   Úitems     r#   ÚputÚSimpleQueue.put%   sK   € à�
‰
×Ñ˜$ÔØ�?‰?Ñ&¨t¯©×/CÑ/C×/EÑ/EØ�O‰O×&Ñ& tÕ,ð 0FÐ&r&   c              ƒ  ó¦  #   • U R                   (       d{  U(       d  [        S5      eU R                  b   S5       eU R                  R	                  5       U l         U R                  I Sh  v•N   U R                  R                  5         SU l        U R                   R                  5       $  N?! U R                  R                  5         SU l        f = f7f)z?Remove and return an item from the queue, waiting if necessary.ústream of frames endedNzcannot call get() concurrently)r    ÚEOFErrorr   r   Úcreate_futureÚcancelÚpopleft)r"   Úblocks     r#   ÚgetÚSimpleQueue.get+   sž   é € à�z�zÞÜÐ7Ó8Ð8Ø—?‘?Ñ*ÐLÐ,LÓLÐ*Ø"Ÿi™i×5Ñ5Ó7ˆDŒOð'Ø—o‘o×%Ð%à—‘×&Ñ&Ô(Ø"&�”Ø�z‰z×!Ñ!Ó#Ð#ñ	 &øà—‘×&Ñ&Ô(Ø"&�•üs0   ‚ACÁB+ Á)B)Á*B+ Á.;CÂ)B+ Â+#CÃCc                ó’   • U R                   b   S5       eU R                  (       a   S5       eU R                  R                  U5        g)z)Put back items into an empty, idle queue.Nz%cannot reset() while get() is runningz&cannot reset() while queue isn't empty)r   r    Úextend)r"   Úitemss     r#   ÚresetÚSimpleQueue.reset9   s<   € à�‰Ñ&ÐOÐ(OÓOÐ&Ø—:—:ÐGÐGÓGˆ~Ø�
‰
×Ñ˜%Õ r&   c                ó¨   • U R                   bE  U R                   R                  5       (       d%  U R                   R                  [        S5      5        ggg)z8Close the queue, raising EOFError in get() if necessary.Nr3   )r   r-   Úset_exceptionr4   r!   s    r#   ÚabortÚSimpleQueue.abort?   s?   € à�?‰?Ñ&¨t¯©×/CÑ/C×/EÑ/EØ�O‰O×)Ñ)¬(Ð3KÓ*LÕMð 0FÐ&r&   )r   r   r    N©ÚreturnÚNone)rE   Úint)r/   r   rE   rF   )T)r8   ÚboolrE   r   )r=   zIterable[T]rE   rF   )Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r$   r)   r0   r9   r>   rB   Ú__static_attributes__© r&   r#   r   r      s&   † ñô?ô
ô-ö$ô!÷Nr&   r   c                  ó  • \ rS rSrSrSSS S 4         SS jjr\SS j5       r\SS j5       r\SSS	 jj5       rSSS
 jjr\SS j5       r\SS j5       r\SSS jj5       rSSS jjrSS jr	SS jr
SS jrSS jrSrg)r   éE   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                 ó   • g r   rO   rO   r&   r#   Ú<lambda>ÚAssembler.<lambda>X   s   € ¨4r&   c                 ó   • g r   rO   rO   r&   r#   rS   rT   Y   s   € ¨Dr&   c                óú   • [        5       U l        Ub  Uc  US-  nUc  Ub  US-  nUb$  Ub!  US:  a  [        S5      eX:  a  [        S5      eXsU l        U l        X0l        X@l        SU l        SU l        SU l	        g )Né   r   z%low must be positive or equal to zeroz)high must be greater than or equal to lowF)
r   ÚframesÚ
ValueErrorÚhighÚlowÚpauseÚresumeÚpausedÚget_in_progressÚclosed)r"   rZ   r[   r\   r]   s        r#   r$   ÚAssembler.__init__T   s”   € ô +6«-ˆŒð Ñ ¡Ø˜!‘)ˆCØ‰<˜C™OØ˜‘7ˆDØÑ ¡Ø�Q‹wÜ Ð!HÓIÐIØ‹zÜ Ð!LÓMÐMØ"ÐˆŒ	�4”8ØŒ
ØŒØˆŒð  %ˆÔð ˆ�r&   c              ƒ  ó   #   • g 7fr   rO   ©r"   Údecodes     r#   r9   ÚAssembler.getv   s   é € Ø7:ùó   ‚c              ƒ  ó   #   • g 7fr   rO   rc   s     r#   r9   re   y   s   é € Ø:=ùrf   c              ƒ  ó   #   • g 7fr   rO   rc   s     r#   r9   re   |   s   é € Ø=@ùrf   c              ƒ  ój  #   • U R                   (       a  [        S5      eSU l          U R                  R                  U R                  (       + 5      I Sh  v•N nU R                  5         UR                  [        L d  UR                  [        L d   eUc  UR                  [        L nU/nUR                  (       d|   U R                  R                  U R                  (       + 5      I Sh  v•N nU R                  5         UR                  [        L d   eUR                  U5        UR                  (       d  M|  SU l         SR                  S U 5       5      nU(       a  UR!                  5       $ U$  GN Nˆ! [        R                   a    U R                  R                  U5        e f = f! SU l         f = f7f)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:
    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 runningTNFr&   c              3  ó8   #   • U  H  oR                   v •  M     g 7fr   )Údata)Ú.0Úframes     r#   Ú	<genexpr>Ú Assembler.get.<locals>.<genexpr>¶   s   é € Ð7² uŸ
ž
²ùs   ‚)r_   r   rX   r9   r`   Úmaybe_resumeÚopcoder   r   Úfinr   ÚCancelledErrorr>   r   r,   Újoinrd   )r"   rd   rn   rX   rl   s        r#   r9   re      s\  é € ð, ××Ü"Ð#KÓLÐLØ#ˆÔð
	)àŸ+™+Ÿ/™/¨d¯k©k¬/Ó:×:ˆEØ×ÑÔØ—<‘<¤7Ò*¨e¯l©l¼iÒ.GÐGÐGØ‰~ØŸ™¬Ð0�Ø�WˆFð —i—iðØ"&§+¡+§/¡/°d·k±k´/Ó"B×B�Eð ×!Ñ!Ô#Ø—|‘|¤wÒ.Ð.Ð.Ø—‘˜eÔ$ð —i—i‘ið $)ˆDÔ ð �x‰xÑ7±Ó7Ó7ˆÞØ—;‘;“=Ð àˆKò9 ;ñ CøÜ×-Ñ-ó ð —K‘K×%Ñ% fÔ-Øð	ûð $)ˆDÕ üsZ   ‚$F3§-F' ÁE.ÁA%F' Â;-E3 Ã(E1Ã)E3 Ã-AF' Ä68F3Å.F' Å1E3 Å31F$Æ$F' Æ'	F0Æ0F3c                ó   • g r   rO   rc   s     r#   Úget_iterÚAssembler.get_iter¼   s   € ØEHr&   c                ó   • g r   rO   rc   s     r#   rw   rx   ¿   s   € ØHKr&   c                ó   • g r   rO   rc   s     r#   rw   rx   Â   s   € ØKNr&   c               óØ  #   • U R                   (       a  [        S5      eSU l          U R                  R                  U R                  (       + 5      I Sh  v•N nU R                  5         UR                  [        L d  UR                  [        L d   eUc  UR                  [        L nU(       a4  [        5       nUR                  UR                  UR                  5      7v •  O[        UR                  5      7v •  UR                  (       d³  U R                  R                  U R                  (       + 5      I Sh  v•N nU R                  5         UR                  [         L d   eU(       a*  WR                  UR                  UR                  5      7v •  O[        UR                  5      7v •  UR                  (       d  M³  SU l         g GNq! [
        R                   a	    SU l         e f = f N°7f)a0  
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.

rj   TNF)r_   r   rX   r9   r`   r   rt   rq   rr   r   r   ÚUTF8Decoderrd   rl   rs   Úbytesr   )r"   rd   rn   Údecoders       r#   rw   rx   Å   sf  é € ð2 ××Ü"Ð#KÓLÐLØ#ˆÔð	ØŸ+™+Ÿ/™/¨d¯k©k¬/Ó:×:ˆEð 	×ÑÔØ�|‰|œwÒ&¨%¯,©,¼)Ò*CÐCÐCØ‰>Ø—\‘\¤WÐ,ˆFÞÜ!“mˆGØ—.‘. §¡¨U¯Y©YÓ7Ô7ô ˜Ÿ
™
Ó#Ó#ð —)—)ð
 Ÿ+™+Ÿ/™/¨d¯k©k¬/Ó:×:ˆEØ×ÑÔØ—<‘<¤7Ò*Ð*Ð*ÞØ—n‘n U§Z¡Z°·±Ó;Ô;ô ˜EŸJ™JÓ'Ó'ð —)—)‘)ð  %ˆÕò= ;øÜ×%Ñ%ó 	Ø#(ˆDÔ Øð	úñ( ;ùsB   ‚$G*§-G ÁGÁG ÁCG*Ä7G(Ä8BG*Æ=G*ÇG ÇG%Ç%G*c                ó’   • U R                   (       a  [        S5      eU R                  R                  U5        U R	                  5         g)z_
Add ``frame`` to the next message.

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

r3   N)r`   r4   rX   r0   Úmaybe_pause)r"   rn   s     r#   r0   ÚAssembler.put
  s3   € ð �;�;ÜÐ3Ó4Ð4à�‰�‰˜ÔØ×ÑÕr&   c                óº   • U R                   c  g[        U R                  5      U R                   :”  a*  U R                  (       d  SU l        U R	                  5         ggg)z7Pause the writer if queue is above the high water mark.NT)rZ   r(   rX   r^   r\   r!   s    r#   r€   ÚAssembler.maybe_pause  sF   € ð �9‰9ÑØô ˆt�{‰{Ó˜dŸi™iÓ'°··ØˆDŒKØ�J‰J�Lð 1<Ð'r&   c                óº   • U R                   c  g[        U R                  5      U R                   ::  a*  U R                  (       a  SU l        U R	                  5         ggg)z7Resume the writer if queue is below the low water mark.NF)r[   r(   rX   r^   r]   r!   s    r#   rq   ÚAssembler.maybe_resume#  sF   € ð �8‰8ÑØô ˆt�{‰{Ó˜tŸx™xÓ'¨D¯K¯KØˆDŒKØ�K‰K�Mð -8Ð'r&   c                ój   • U R                   (       a  gSU l         U R                  R                  5         g)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`.

NT)r`   rX   rB   r!   s    r#   ÚcloseÚAssembler.close.  s'   € ð �;�;ØàˆŒð 	�‰×ÑÕr&   )r`   rX   r_   rZ   r[   r\   r^   r]   )
rZ   ú
int | Noner[   r‰   r\   úCallable[[], Any]r]   rŠ   rE   rF   )rd   úLiteral[True]rE   Ústr)rd   úLiteral[False]rE   r}   r   )rd   úbool | NonerE   r   )rd   r‹   rE   zAsyncIterator[str])rd   r�   rE   zAsyncIterator[bytes])rd   rŽ   rE   zAsyncIterator[Data])rn   r   rE   rF   rD   )rI   rJ   rK   rL   rM   r$   r   r9   rw   r0   r€   rq   r‡   rN   rO   r&   r#   r   r   E   sÄ   † ñð   ØÙ#/Ù$0ð àð ð ð ð !ð	 ð
 "ð ð 
õ ðD Û:ó Ø:àÛ=ó Ø=àÝ@ó Ø@ö;ðz ÛHó ØHàÛKó ØKàÝNó ØNöC%ôJô	ô	÷r&   )Ú
__future__r   r   Úcodecsr   Úcollections.abcr   r   Útypingr   r   r   r	   r
   r   Ú
exceptionsr   rX   r   r   r   r   r   Ú__all__Úgetincrementaldecoderr|   r   r   r   rO   r&   r#   Ú<module>r–      sg   ðÝ "ã Û Û ß 3ß E× Eå )ß 7Ó 7Ý ð ˆ-€à×*Ò*¨7Ó3€áˆCƒL€ô-N�'˜!‘*ô -N÷`wò wr&   