o
    Ð­jy+  ã                   @  sÂ   d dl mZ d dlZd dlZd dlZd dlmZm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	gZe d
¡ZedƒZG dd„ dee ƒZG dd	„ d	ƒZdS )é    )ÚannotationsN)ÚAsyncIteratorÚIterable)ÚAnyÚCallableÚGenericÚLiteralÚTypeVarÚoverloadé   )ÚConcurrencyError)Ú	OP_BINARYÚOP_CONTÚOP_TEXTÚFrame)ÚDataÚ	Assemblerzutf-8ÚTc                   @  sN   e Zd ZdZddd„Zddd„Zddd„Zdddd„Zddd„Zddd„Z	dS )ÚSimpleQueuez…
    Simplified version of :class:`asyncio.Queue`.

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

    ÚreturnÚNonec                 C  s   t  ¡ | _d | _t ¡ | _d S ©N)ÚasyncioÚget_running_loopÚloopÚ
get_waiterÚcollectionsÚdequeÚqueue©Úself© r!   úX/var/www/html/CropPilot/venv/lib/python3.10/site-packages/websockets/asyncio/messages.pyÚ__init__   s   
zSimpleQueue.__init__Úintc                 C  s
   t | jƒS r   )Úlenr   r   r!   r!   r"   Ú__len__"   s   
zSimpleQueue.__len__Úitemr   c                 C  s8   | j  |¡ | jdur| j ¡ s| j d¡ dS dS dS )zPut an item into the queue.N)r   Úappendr   ÚdoneÚ
set_result)r    r'   r!   r!   r"   Úput%   s   ÿzSimpleQueue.putTÚblockÚboolc                 Ã  sp   �| j s3|s
tdƒ‚| jdu sJ dƒ‚| j ¡ | _z| jI dH  W | j ¡  d| _n	| j ¡  d| _w | j  ¡ S )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    r,   r!   r!   r"   Úget+   s   €

ÿ
zSimpleQueue.getÚitemsúIterable[T]c                 C  s0   | j du s	J dƒ‚| jrJ dƒ‚| j |¡ dS )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    r4   r!   r!   r"   Úreset9   s   zSimpleQueue.resetc                 C  s0   | j dur| j  ¡ s| j  tdƒ¡ dS dS dS )z8Close the queue, raising EOFError in get() if necessary.Nr.   )r   r)   Úset_exceptionr/   r   r!   r!   r"   Úabort?   s   ÿzSimpleQueue.abortN©r   r   )r   r$   )r'   r   r   r   )T)r,   r-   r   r   )r4   r5   r   r   )
Ú__name__Ú
__module__Ú__qualname__Ú__doc__r#   r&   r+   r3   r7   r9   r!   r!   r!   r"   r      s    



r   c                   @  sÄ   e Zd ZdZdddd„ dd„ fd.dd„Zed/dd„ƒZed0dd„ƒZed1d2dd„ƒZd1d2dd„Zed3dd„ƒZed4d d„ƒZed1d5d"d„ƒZd1d5d#d„Zd6d&d'„Zd7d(d)„Z	d7d*d+„Z
d7d,d-„ZdS )8r   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                   C  ó   d S r   r!   r!   r!   r!   r"   Ú<lambda>X   ó    zAssembler.<lambda>c                   C  r?   r   r!   r!   r!   r!   r"   r@   Y   rA   Úhighú
int | NoneÚlowÚpauseúCallable[[], Any]Úresumer   r   c                 C  s˜   t ƒ | _|d ur|d u r|d }|d u r|d ur|d }|d ur4|d ur4|dk r,tdƒ‚||k r4tdƒ‚||| _| _|| _|| _d| _d| _d| _	d S )Né   r   z%low must be positive or equal to zeroz)high must be greater than or equal to lowF)
r   ÚframesÚ
ValueErrorrB   rD   rE   rG   ÚpausedÚget_in_progressÚclosed)r    rB   rD   rE   rG   r!   r!   r"   r#   T   s    
zAssembler.__init__ÚdecodeúLiteral[True]Ústrc                 Ã  ó   �d S r   r!   ©r    rN   r!   r!   r"   r3   v   ó   €zAssembler.getúLiteral[False]Úbytesc                 Ã  rQ   r   r!   rR   r!   r!   r"   r3   y   rS   úbool | Noner   c                 Ã  rQ   r   r!   rR   r!   r!   r"   r3   |   rS   c                 Ã  s  �| j rtdƒ‚d| _ z_| j | j ¡I dH }|  ¡  |jtu s'|jtu s'J ‚|du r0|jtu }|g}|j	sfz| j | j ¡I dH }W n t
jyR   | j |¡ ‚ w |  ¡  |jtu s^J ‚| |¡ |j	r6W d| _ nd| _ w d dd„ |D ƒ¡}|r| ¡ S |S )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 runningTNFó    c                 s  s   � | ]}|j V  qd S r   )Údata)Ú.0Úframer!   r!   r"   Ú	<genexpr>¶   s   € z Assembler.get.<locals>.<genexpr>)rL   r   rI   r3   rM   Úmaybe_resumeÚopcoder   r   Úfinr   ÚCancelledErrorr7   r   r(   ÚjoinrN   )r    rN   r[   rI   rY   r!   r!   r"   r3      s8   €
ü
ö€úAsyncIterator[str]c                 C  r?   r   r!   rR   r!   r!   r"   Úget_iter¼   ó   zAssembler.get_iterúAsyncIterator[bytes]c                 C  r?   r   r!   rR   r!   r!   r"   rc   ¿   rd   úAsyncIterator[Data]c                 C  r?   r   r!   rR   r!   r!   r"   rc   Â   rd   c                 C s  �| j rtdƒ‚d| _ z| j | j ¡I dH }W n tjy$   d| _ ‚ w |  ¡  |jt	u s5|jt
u s5J ‚|du r>|jt	u }|rMtƒ }| |j|j¡V  nt|jƒV  |js�| j | j ¡I dH }|  ¡  |jtu slJ ‚|rx| |j|j¡V  nt|jƒV  |jrVd| _ dS )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.

        rW   TNF)rL   r   rI   r3   rM   r   r`   r]   r^   r   r   ÚUTF8DecoderrN   rY   r_   rU   r   )r    rN   r[   Údecoderr!   r!   r"   rc   Å   s6   €	þ
ô
r[   r   c                 C  s&   | j rtdƒ‚| j |¡ |  ¡  dS )z
        Add ``frame`` to the next message.

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

        r.   N)rM   r/   rI   r+   Úmaybe_pause)r    r[   r!   r!   r"   r+   
  s   zAssembler.putc                 C  s>   | j du rdS t| jƒ| j kr| jsd| _|  ¡  dS dS dS )z7Pause the writer if queue is above the high water mark.NT)rB   r%   rI   rK   rE   r   r!   r!   r"   ri     ó   
þzAssembler.maybe_pausec                 C  s>   | j du rdS t| jƒ| j kr| jrd| _|  ¡  dS dS dS )z7Resume the writer if queue is below the low water mark.NF)rD   r%   rI   rK   rG   r   r!   r!   r"   r]   #  rj   zAssembler.maybe_resumec                 C  s   | j rdS d| _ | j ¡  dS )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)rM   rI   r9   r   r!   r!   r"   Úclose.  s   zAssembler.close)
rB   rC   rD   rC   rE   rF   rG   rF   r   r   )rN   rO   r   rP   )rN   rT   r   rU   r   )rN   rV   r   r   )rN   rO   r   rb   )rN   rT   r   re   )rN   rV   r   rf   )r[   r   r   r   r:   )r;   r<   r=   r>   r#   r
   r3   rc   r+   ri   r]   rk   r!   r!   r!   r"   r   E   s2    û"=
E

)Ú
__future__r   r   Úcodecsr   Úcollections.abcr   r   Útypingr   r   r   r   r	   r
   Ú
exceptionsr   rI   r   r   r   r   r   Ú__all__Úgetincrementaldecoderrg   r   r   r   r!   r!   r!   r"   Ú<module>   s     
0