o
    ØRWjÚ  ã                   @  s:   d dl mZ d dlZd dlZddlmZ G dd„ dƒZdS )é    )ÚannotationsNé   )ÚWebSocketQueueFullErrorc                   @  sX   e Zd ZdZdddd„Zddd„Zd dd„Zd!dd„Zd"dd„Zd#dd„Z	d$dd„Z
dS )%Ú	SendQueuezèBounded byte-size queue for outgoing WebSocket messages.

    Messages are stored as pre-serialized strings. The queue enforces a
    maximum byte budget so that unbounded buffering cannot occur during
    reconnection windows.
    é   Ú	max_bytesÚintÚreturnÚNonec                 C  s    g | _ d| _|| _t ¡ | _d S ©Nr   )Ú_queueÚ_bytesÚ
_max_bytesÚ	threadingÚLockÚ_lock)Úselfr   © r   úa/home/esfera/Documents/content_generation/venv/lib/python3.10/site-packages/openai/_send_queue.pyÚ__init__   s   zSendQueue.__init__ÚdataÚstrc                 C  sp   t | d¡ƒ}| j�$ | j| | jkrtdƒ‚| j ||f¡ |  j|7  _W d  ƒ dS 1 s1w   Y  dS )zŽAppend *data* to the queue.

        Raises :class:`WebSocketQueueFullError` if the message would
        exceed the byte-size limit.
        zutf-8z%send queue is full, message discardedN)ÚlenÚencoder   r   r   r   r   Úappend)r   r   Úbyte_lengthr   r   r   Úenqueue   s   "üzSendQueue.enqueueÚsendútyping.Callable[[str], object]c                 C  sÊ   | j � t| jƒ}| j ¡  d| _W d  ƒ n1 sw   Y  t|ƒD ]>\}\}}z||ƒ W q$ tyb   | j � ||d… }|| j | _tdd„ | jD ƒƒ| _W d  ƒ ‚ 1 s\w   Y  ‚ w dS )z«Send every queued message via *send*.

        If *send* raises, the failing message and all subsequent messages
        are re-queued and the error is re-raised.
        r   Nc                 s  ó   � | ]\}}|V  qd S ©Nr   ©Ú.0Ú_Úblr   r   r   Ú	<genexpr>8   ó   € z'SendQueue.flush_sync.<locals>.<genexpr>©r   Úlistr   Úclearr   Ú	enumerateÚ	ExceptionÚsum©r   r   ÚpendingÚir   Ú_byte_lengthÚ	remainingr   r   r   Ú
flush_sync&   s&   

ý
ýüûýzSendQueue.flush_syncú0typing.Callable[[str], typing.Awaitable[object]]c                 Ã  sÒ   �| j � t| jƒ}| j ¡  d| _W d  ƒ n1 sw   Y  t|ƒD ]A\}\}}z	||ƒI dH  W q% tyf   | j � ||d… }|| j | _tdd„ | jD ƒƒ| _W d  ƒ ‚ 1 s`w   Y  ‚ w dS )z$Async variant of :meth:`flush_sync`.r   Nc                 s  r   r    r   r!   r   r   r   r%   I   r&   z(SendQueue.flush_async.<locals>.<genexpr>r'   r-   r   r   r   Úflush_async;   s(   €

ý
ýüûýzSendQueue.flush_asyncú	list[str]c                 C  sN   | j � dd„ | jD ƒ}| j ¡  d| _|W  d  ƒ S 1 s w   Y  dS )z&Remove and return all queued messages.c                 S  s   g | ]\}}|‘qS r   r   )r"   r   r#   r   r   r   Ú
<listcomp>O   s    z#SendQueue.drain.<locals>.<listcomp>r   N)r   r   r)   r   )r   Úitemsr   r   r   ÚdrainL   s   
$üzSendQueue.drainc                 C  s4   | j � t| jƒW  d   ƒ S 1 sw   Y  d S r    ©r   r   r   ©r   r   r   r   Ú__len__T   s   $ÿzSendQueue.__len__Úboolc                 C  s8   | j � t| jƒdkW  d   ƒ S 1 sw   Y  d S r   r9   r:   r   r   r   Ú__bool__X   s   $ÿzSendQueue.__bool__N)r   )r   r   r	   r
   )r   r   r	   r
   )r   r   r	   r
   )r   r3   r	   r
   )r	   r5   )r	   r   )r	   r<   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r2   r4   r8   r;   r=   r   r   r   r   r      s    




r   )Ú
__future__r   Útypingr   Ú_exceptionsr   r   r   r   r   r   Ú<module>   s
   