
    Yj                    T    d dl mZ d dlZd dlZd dlmZ d dlZddlm	Z	  G d d      Z
y)    )annotationsN)deque   )WebSocketQueueFullErrorc                  Z    e Zd ZdZdddZddZddZddZddZddZ	ddZ
dd	Zdd
Zy)	SendQueuezBounded 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.
    c                n    g | _         d| _        || _        t        j                         | _        d | _        y Nr   )_queue_bytes
_max_bytes	threadingLock_lock_flush_done)self	max_bytess     Q/var/www/html/vmm-quimica/venv/lib/python3.12/site-packages/openai/_send_queue.py__init__zSendQueue.__init__   s-    -/#^^%
37    c                ,   t        |j                  d            }| j                  5  | j                  |z   | j                  kD  rt        d      | j                  j                  ||f       | xj                  |z  c_        ddd       y# 1 sw Y   yxY w)zAppend *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   databyte_lengths      r   enqueuezSendQueue.enqueue   sy     $++g./ZZ 	'{{[(4??:-.UVVKKk23KK;&K		' 	' 	's   AB

Bc                   t        | j                         x}t        j                        r;|j	                          t        | j                         x}t        j                        r;	 |rM|d   \  }} ||       | j
                  5  |j                          | xj                  |z  c_        ddd       |rM| j                  |       y# 1 sw Y   xY w# | j                  |       w xY w)zSend 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   N)	
isinstance_begin_flushr   Eventwaitr   popleftr   
_end_flushr   sendpendingr   r   s        r   
flush_synczSendQueue.flush_sync(   s     D$5$5$77ILLN D$5$5$77I	%$+AJ!kT
ZZ /OO%KK;.K/  OOG$	/ / OOG$s$   'C &C+
C CC C'c                :  K   t        | j                         x}t        j                        r^t        j
                  j                  |j                  d       d{    t        | j                         x}t        j                        r^	 |rU|d   \  }} ||       d{    | j                  5  |j                          | xj                  |z  c_
        ddd       |rU| j                  |       y7 7 U# 1 sw Y   "xY w# | j                  |       w xY ww)z$Async variant of :meth:`flush_sync`.T)abandon_on_cancelNr   )r   r    r   r!   anyio	to_threadrun_syncr"   r   r#   r   r$   r%   s        r   flush_asynczSendQueue.flush_async;   s     D$5$5$77I //**7<<4*PPP D$5$5$77I
	%$+AJ!k4j  ZZ /OO%KK;.K/  OOG$ Q
 !/ / OOG$sT   ADC5.DD !C7"D 2&C9
D #D7D 9D>D DDc                   | j                   5  | j                  | j                  cd d d        S t        | j                        }| j                  j	                          t        j                         | _        |cd d d        S # 1 sw Y   y xY wN)r   r   r   r   clearr   r!   r   r'   s     r   r    zSendQueue._begin_flushL   so    ZZ 	+''	 	 DKK(GKK(0D	 	 	s   BA	BBc                    | j                   5  t        |      | j                  z   | _        | j                  J | j                  j	                          d | _        d d d        y # 1 sw Y   y xY wr0   )r   listr   r   setr2   s     r   r$   zSendQueue._end_flushV   s^    ZZ 	$w-$++5DK##///  "#D		$ 	$ 	$s   AA##A,c                $   | j                   5  | j                  D cg c]  \  }}|	 }}}| xj                  t        d | j                  D              z  c_        | j                  j	                          |cddd       S c c}}w # 1 sw Y   yxY w)z&Remove and return all queued messages.c              3  &   K   | ]	  \  }}|  y wr0    ).0_r   s      r   	<genexpr>z"SendQueue.drain.<locals>.<genexpr>a   s     M~q+{Ms   N)r   r   r   sumr1   )r   r   r:   itemss       r   drainzSendQueue.drain]   sq    ZZ 	)-5gdAT5E5KK3MMMMKKK		 	5	 	s   BB AB BBc                p    | j                   5  t        | j                        cd d d        S # 1 sw Y   y xY wr0   r   r   r   r   s    r   __len__zSendQueue.__len__e   s*    ZZ 	$t{{#	$ 	$ 	$s   ,5c                v    | j                   5  t        | j                        dkD  cd d d        S # 1 sw Y   y xY wr
   r@   rA   s    r   __bool__zSendQueue.__bool__i   s/    ZZ 	(t{{#a'	( 	( 	(s   /8N)i   )r   intreturnNone)r   strrF   rG   )r&   ztyping.Callable[[str], object]rF   rG   )r&   z0typing.Callable[[str], typing.Awaitable[object]]rF   rG   )rF   z(deque[tuple[str, int]] | threading.Event)r'   zdeque[tuple[str, int]]rF   rG   )rF   z	list[str])rF   rE   )rF   bool)__name__
__module____qualname____doc__r   r   r(   r.   r    r$   r>   rB   rD   r8   r   r   r   r      s4    8'%&%"$$(r   r   )
__future__r   typingr   collectionsr   anyio.to_threadr+   _exceptionsr   r   r8   r   r   <module>rS      s#    "     0_( _(r   