o
    Ð­j:2  ã                   @  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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 d¡ZG dd
„ d
ƒZdS )é    )ÚannotationsN)ÚAnyÚCallableÚIterableÚIteratorÚLiteralÚoverloadé   )ÚConcurrencyError)Ú	OP_BINARYÚOP_CONTÚOP_TEXTÚFrame)ÚDataé   )ÚDeadlineÚ	Assemblerzutf-8c                   @  sú   e Zd ZdZdddd„ dd„ fd8dd„Zd9d:dd„Zd;dd„Zed<dd„ƒZed=d d„ƒZed9d<d!d„ƒZed9d=d"d„ƒZed>d?d%d„ƒZd>d?d&d„Zed@d(d)„ƒZ	edAd+d)„ƒZ	ed9dBd-d)„ƒZ	d9dBd.d)„Z	dCd0d1„Z
dDd2d3„ZdDd4d5„ZdDd6d7„ZdS )Er   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 ©N© r   r   r   úU/var/www/html/CropPilot/venv/lib/python3.10/site-packages/websockets/sync/messages.pyÚ<lambda>&   ó    zAssembler.<lambda>c                   C  r   r   r   r   r   r   r   r   '   r   Úhighú
int | NoneÚlowÚpauseúCallable[[], Any]ÚresumeÚreturnÚNonec                 C  s¤   t  ¡ | _t ¡ | _|d ur|d u r|d }|d u r"|d ur"|d }|d ur:|d ur:|dk r2tdƒ‚||k r:t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)Ú	threadingÚLockÚmutexÚqueueÚSimpleQueueÚframesÚ
ValueErrorr   r   r   r   ÚpausedÚget_in_progressÚclosed)Úselfr   r   r   r   r   r   r   Ú__init__"   s"   
	

zAssembler.__init__Útimeoutúfloat | Noner   c                 C  s¢   | j rz	| jjdd�}W n: tjy   tdƒd ‚w z|d ur+|dkr+| jjdd�}n| jjd|d�}W n tjyF   td|d›d	�ƒd ‚w |d u rOtdƒ‚|S )
NF©Úblockústream of frames endedr   T)r1   r.   ztimed out in z.1fÚs)r+   r'   Úgetr%   ÚEmptyÚEOFErrorÚTimeoutError)r,   r.   Úframer   r   r   Úget_next_frameH   s"   
ÿ€ÿzAssembler.get_next_framer'   úIterable[Frame]c              	   C  sŠ   | j �8 g }z	 | | jjdd�¡ q tjy   Y nw |D ]}| j |¡ q|D ]}| j |¡ q*W d   ƒ d S 1 s>w   Y  d S )NTFr0   )r$   Úappendr'   r4   r%   r5   Úput)r,   r'   Úqueuedr8   r   r   r   Úreset_queue^   s   ÿÿÿ"özAssembler.reset_queueÚdecodeúLiteral[True]Ústrc                 C  r   r   r   ©r,   r.   r?   r   r   r   r4   t   ó   zAssembler.getúLiteral[False]Úbytesc                 C  r   r   r   rB   r   r   r   r4   w   rC   c                C  r   r   r   rB   r   r   r   r4   z   rC   c                C  r   r   r   rB   r   r   r   r4   }   rC   úbool | Noner   c                 C  r   r   r   rB   r   r   r   r4   €   rC   c                 C  sn  | j � | jrtdƒ‚d| _W d  ƒ n1 sw   Y  zƒt|ƒ}|  |jdd�¡}| j � |  ¡  W d  ƒ n1 s=w   Y  |jtu sN|jt	u sNJ ‚|du rW|jtu }|g}|j
sœz|  |jdd�¡}W n tyu   |  |¡ ‚ w | j � |  ¡  W d  ƒ n1 sˆw   Y  |jtu s”J ‚| |¡ |j
r]W d| _nd| _w d dd„ |D ƒ¡}|rµ| ¡ S |S )	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_elapsedó    c                 s  s   � | ]}|j V  qd S r   )Údata)Ú.0r8   r   r   r   Ú	<genexpr>Å   s   € z Assembler.get.<locals>.<genexpr>)r$   r*   r
   r   r9   r.   Úmaybe_resumeÚopcoder   r   Úfinr7   r>   r   r;   Újoinr?   )r,   r.   r?   Údeadliner8   r'   rJ   r   r   r   r4   ƒ   sH   ý
ÿ

ÿ
ü
ÿ
ó€úIterator[str]c                 C  r   r   r   ©r,   r?   r   r   r   Úget_iterË   rC   zAssembler.get_iterúIterator[bytes]c                 C  r   r   r   rS   r   r   r   rT   Î   rC   úIterator[Data]c                 C  r   r   r   rS   r   r   r   rT   Ñ   rC   c                 c  sD  � | j � | jrtdƒ‚d| _W d  ƒ n1 sw   Y  |  ¡ }| j � |  ¡  W d  ƒ n1 s4w   Y  |jtu sE|jtu sEJ ‚|du rN|jtu }|r]tƒ }| 	|j
|j¡V  nt|j
ƒV  |js�|  ¡ }| j � |  ¡  W d  ƒ n1 s|w   Y  |jtu sˆJ ‚|r”| 	|j
|j¡V  nt|j
ƒV  |jrfd| _dS )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.

        rG   TNF)r$   r*   r
   r9   rM   rN   r   r   ÚUTF8Decoderr?   rJ   rO   rE   r   )r,   r?   r8   Údecoderr   r   r   rT   Ô   s8   €ý
ÿ

ÿ÷
r8   c                 C  sN   | j � | jrtdƒ‚| j |¡ |  ¡  W d  ƒ dS 1 s w   Y  dS )z
        Add ``frame`` to the next message.

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

        r2   N)r$   r+   r6   r'   r<   Úmaybe_pause)r,   r8   r   r   r   r<     s   
"ûzAssembler.putc                 C  sL   | j du rdS | j ¡ sJ ‚| 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)r   r$   Úlockedr'   Úqsizer)   r   ©r,   r   r   r   rY   -  ó   
þzAssembler.maybe_pausec                 C  sL   | j du rdS | j ¡ sJ ‚| 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)r   r$   rZ   r'   r[   r)   r   r\   r   r   r   rM   :  r]   zAssembler.maybe_resumec                 C  s€   | j �3 | jr	 W d  ƒ dS d| _| jr| j d¡ | jr.d| _|  ¡  W d  ƒ dS W d  ƒ dS 1 s9w   Y  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`.

        NTF)r$   r+   r*   r'   r<   r)   r   r\   r   r   r   ÚcloseG  s   þ
ó
"özAssembler.close)
r   r   r   r   r   r   r   r   r   r    r   )r.   r/   r   r   )r'   r:   r   r    )r.   r/   r?   r@   r   rA   )r.   r/   r?   rD   r   rE   )NN)r.   r/   r?   rF   r   r   )r?   r@   r   rR   )r?   rD   r   rU   )r?   rF   r   rV   )r8   r   r   r    )r   r    )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r-   r9   r>   r   r4   rT   r<   rY   rM   r^   r   r   r   r   r      s>    û&
H
@

)Ú
__future__r   Úcodecsr%   r"   Útypingr   r   r   r   r   r   Ú
exceptionsr
   r'   r   r   r   r   r   Úutilsr   Ú__all__ÚgetincrementaldecoderrW   r   r   r   r   r   Ú<module>   s     
