Ë
    ÚŠ±jì^  ã                   óD  — U d Z ddlZddlZddlZddlZddlZddlZddlZddlm	Z	m
Z
mZ ddlmZmZ  G d„ d«      Z G d„ d	«      Z G d
„ d«      Z ej$                  «       Z ej(                  d«      Zej(                  e   ed<   dad„ Zd„ Zddefd„Zddefd„Zd„ Zd„ Zy)zu
Dirty Client

Client for HTTP workers to communicate with the dirty worker pool.
Provides both sync and async APIs.
é    Né   )ÚDirtyConnectionErrorÚ
DirtyErrorÚDirtyTimeoutError)ÚDirtyProtocolÚmake_requestc                   óx   — e Zd ZdZdd„Zd„ Zd„ Zd„ Zd„ Zd„ Z	d„ Z
d	„ Zd
„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zy)ÚDirtyClientzä
    Client for calling dirty workers from HTTP workers.

    Provides both sync and async APIs. The sync API is for traditional
    sync workers (sync, gthread), while the async API is for async
    workers (asgi, gevent).
    c                 ó|   — || _         || _        d| _        d| _        d| _        t        j                  «       | _        y)z½
        Initialize the dirty client.

        Args:
            socket_path: Path to the dirty arbiter's Unix socket
            timeout: Default timeout for operations in seconds
        N)Úsocket_pathÚtimeoutÚ_sockÚ_readerÚ_writerÚ	threadingÚLockÚ_lock)Úselfr   r   s      úZ/var/www/io.vulcan-creative.com/venv/lib/python3.12/site-packages/gunicorn/dirty/client.pyÚ__init__zDirtyClient.__init__(   s6   € ð 'ˆÔØˆŒØˆŒ
ØˆŒØˆŒÜ—^‘^Ó%ˆ�
ó    c                 ó¨  — | j                   �y	 t        j                  t        j                  t        j                  «      | _         | j                   j	                  | j
                  «       | j                   j                  | j                  «       y# t        j                  t        f$ r'}d| _         t        d|› �| j                  ¬«      |‚d}~ww xY w)z…
        Establish sync socket connection to arbiter.

        Raises:
            DirtyConnectionError: If connection fails
        Nú$Failed to connect to dirty arbiter: ©r   )r   ÚsocketÚAF_UNIXÚSOCK_STREAMÚ
settimeoutr   Úconnectr   ÚerrorÚOSErrorr   ©r   Úes     r   r   zDirtyClient.connect;   s¥   € ð �:‰:Ð!Øð		ÜŸ™¤v§~¡~´v×7IÑ7IÓJˆDŒJØ�J‰J×!Ñ! $§,¡,Ô/Ø�J‰J×Ñ˜t×/Ñ/Õ0øÜ—‘œgÐ&ò 	ØˆDŒJÜ&Ø6°q°cÐ:Ø ×,Ñ,ôð ðûð	ús   �BB ÂCÂ*"CÃCc                 ón   — | j                   5  | j                  ||||«      cddd«       S # 1 sw Y   yxY w)a  
        Execute an action on a dirty app (sync/blocking).

        Args:
            app_path: Import path of the dirty app (e.g., 'myapp.ml:MLApp')
            action: Action to call on the app
            *args: Positional arguments
            **kwargs: Keyword arguments

        Returns:
            Result from the dirty app action

        Raises:
            DirtyConnectionError: If connection fails
            DirtyTimeoutError: If operation times out
            DirtyError: If execution fails
        N)r   Ú_execute_locked©r   Úapp_pathÚactionÚargsÚkwargss        r   ÚexecutezDirtyClient.executeP   s8   € ð$ �Z‰Zñ 	HØ×'Ñ'¨°&¸$ÀÓG÷	H÷ 	Hò 	Hús   �+«4c                 ó*  — | j                   €| j                  «        t        t        j                  «       «      }t        |||||¬«      }	 t        j                  | j                   |«       t        j                  | j                   «      }| j                  |«      S # t        j                  $ r( | j                  «        t        d| j                  ¬«      ‚t        $ r5}| j                  «        t        |t         «      r‚ t#        d|› �«      |‚d}~ww xY w)zExecute while holding the lock.N©Ú
request_idr'   r(   r)   r*   ú&Timeout waiting for dirty app response©r   úCommunication error: )r   r   ÚstrÚuuidÚuuid4r   r   Úwrite_messageÚread_messageÚ_handle_responser   r   Ú_close_socketr   Ú	ExceptionÚ
isinstancer   r   ©	r   r'   r(   r)   r*   r.   ÚrequestÚresponser#   s	            r   r%   zDirtyClient._execute_lockede   sö   € ð �:‰:ÐØ�L‰LŒNô œŸ™›Ó&ˆ
ÜØ!ØØØØô
ˆð	Kä×'Ñ'¨¯
©
°GÔ<ô %×1Ñ1°$·*±*Ó=ˆHð ×(Ñ(¨Ó2Ð2øÜ�~‰~ò 	Ø×ÑÔ Ü#Ø8ØŸ™ôð ô ò 	KØ×ÑÔ Ü˜!œZÔ(ØÜ&Ð)>¸q¸cÐ'BÓCÈÐJûð		Kús   ÁAB ÂADÃ0DÄDc                 ó    — t        | ||||«      S )a*  
        Stream results from a dirty app action (sync).

        This method returns an iterator that yields chunks from a streaming
        response. Use this for actions that return generators.

        Args:
            app_path: Import path of the dirty app (e.g., 'myapp.ml:MLApp')
            action: Action to call on the app
            *args: Positional arguments
            **kwargs: Keyword arguments

        Yields:
            Chunks of data from the streaming response

        Raises:
            DirtyConnectionError: If connection fails
            DirtyTimeoutError: If operation times out
            DirtyError: If execution fails

        Example::

            for chunk in client.stream("myapp.llm:LLMApp", "generate", prompt):
                print(chunk, end="", flush=True)
        )ÚDirtyStreamIteratorr&   s        r   ÚstreamzDirtyClient.streamŠ   s   € ô4 # 4¨°6¸4ÀÓHÐHr   c                 ó   — |j                  d«      }|t        j                  k(  r|j                  d«      S |t        j                  k(  r)|j                  di «      }t	        j
                  |«      }|‚t	        d|› �«      ‚)z<Handle response message, extracting result or raising error.ÚtypeÚresultr    zUnknown response type: )Úgetr   ÚMSG_TYPE_RESPONSEÚMSG_TYPE_ERRORr   Ú	from_dict)r   r=   Úmsg_typeÚ
error_infor    s        r   r7   zDirtyClient._handle_response¦   ss   € à—<‘< Ó'ˆà”}×6Ñ6Ò6Ø—<‘< Ó)Ð)Øœ×5Ñ5Ò5Ø!Ÿ™ g¨rÓ2ˆJÜ×(Ñ(¨Ó4ˆEØˆKäÐ6°x°jÐAÓBÐBr   c                 ó€   — | j                   �#	 | j                   j                  «        d| _         yy# t        $ r Y Œw xY w)zClose the socket connection.N)r   Úcloser9   ©r   s    r   r8   zDirtyClient._close_socket³   sC   € à�:‰:Ð!ðØ—
‘
× Ñ Ô"ð ˆD�Jð "øô ò Ùðús   Ž1 ±	=¼=c                 óf   — | j                   5  | j                  «        ddd«       y# 1 sw Y   yxY w)zClose the sync connection.N)r   r8   rL   s    r   rK   zDirtyClient.close¼   s*   € à�Z‰Zñ 	!Ø×ÑÔ ÷	!÷ 	!ñ 	!ús   �'§0c              ƒ   óˆ  K  — | j                   �y	 t        j                  t        j                  | j                  «      | j
                  ¬«      ƒ d{  –—† \  | _        | _         y7 Œ# t        j                  $ r t        d| j
                  ¬«      ‚t        t        f$ r }t        d|› �| j                  ¬«      |‚d}~ww xY w­w)z
        Establish async connection to arbiter.

        Raises:
            DirtyConnectionError: If connection fails
        Nr0   z#Timeout connecting to dirty arbiterr   r   )r   ÚasyncioÚwait_forÚopen_unix_connectionr   r   r   ÚTimeoutErrorr   r!   ÚConnectionErrorr   r"   s     r   Úconnect_asynczDirtyClient.connect_asyncÅ   sº   è ø€ ð �<‰<Ð#Øð	Ü/6×/?Ñ/?Ü×,Ñ,¨T×-=Ñ-=Ó>ØŸ™ô0÷ *Ñ&ˆDŒL˜$�,ð *ùô ×#Ñ#ò 	Ü#Ø5ØŸ™ôð ô œÐ)ò 	Ü&Ø6°q°cÐ:Ø ×,Ñ,ôð ðûð	üs;   ‚C‘AA' ÁA%ÁA' Á$CÁ%A' Á'8B?ÂB:Â:B?Â?Cc              �   óÐ  K  — | j                   €| j                  «       ƒ d{  –—†  t        t        j                  «       «      }t        |||||¬«      }	 t        j                  | j                   |«      ƒ d{  –—†  t        j                  t        j                  | j                  «      | j                  ¬«      ƒ d{  –—† }| j                  |«      S 7 Œ±7 Œ]7 Œ# t        j                  $ r1 | j                  «       ƒ d{  –—†7   t!        d| j                  ¬«      ‚t"        $ r>}| j                  «       ƒ d{  –—†7   t%        |t&        «      r‚ t)        d|› �«      |‚d}~ww xY w­w)aï  
        Execute an action on a dirty app (async/non-blocking).

        Args:
            app_path: Import path of the dirty app
            action: Action to call on the app
            *args: Positional arguments
            **kwargs: Keyword arguments

        Returns:
            Result from the dirty app action

        Raises:
            DirtyConnectionError: If connection fails
            DirtyTimeoutError: If operation times out
            DirtyError: If execution fails
        Nr-   r0   r/   r1   )r   rT   r2   r3   r4   r   r   Úwrite_message_asyncrO   rP   Úread_message_asyncr   r   r7   rR   Ú_close_asyncr   r9   r:   r   r   r;   s	            r   Úexecute_asynczDirtyClient.execute_asyncß   sD  è ø€ ð& �<‰<ÐØ×$Ñ$Ó&×&Ð&ô œŸ™›Ó&ˆ
ÜØ!ØØØØô
ˆð	Kä×3Ñ3°D·L±LÀ'ÓJ×JÐJô %×-Ñ-Ü×0Ñ0°·±Ó>ØŸ™ô÷ ˆHð ×(Ñ(¨Ó2Ð2ð/ 'øð Køðùô ×#Ñ#ò 	Ø×#Ñ#Ó%×%Ñ%Ü#Ø8ØŸ™ôð ô ò 	KØ×#Ñ#Ó%×%Ñ%Ü˜!œZÔ(ØÜ&Ð)>¸q¸cÐ'BÓCÈÐJûð		Küsp   ‚ E&¢C£1E&Á#C Á8CÁ9AC Â>CÂ?C ÃE&ÃC ÃC Ã&E#Ä DÄ$E#Ä%EÄ8D;Ä9%EÅE#Å#E&c                 ó    — t        | ||||«      S )a8  
        Stream results from a dirty app action (async).

        This method returns an async iterator that yields chunks from a
        streaming response. Use this for actions that return generators.

        Args:
            app_path: Import path of the dirty app (e.g., 'myapp.ml:MLApp')
            action: Action to call on the app
            *args: Positional arguments
            **kwargs: Keyword arguments

        Yields:
            Chunks of data from the streaming response

        Raises:
            DirtyConnectionError: If connection fails
            DirtyTimeoutError: If operation times out
            DirtyError: If execution fails

        Example::

            async for chunk in client.stream_async("myapp.llm:LLMApp", "generate", prompt):
                await response.write(chunk)
        )ÚDirtyAsyncStreamIteratorr&   s        r   Ústream_asynczDirtyClient.stream_async  s   € ô4 (¨¨h¸ÀÀfÓMÐMr   c              ƒ   óÞ   K  — | j                   �L	 | j                   j                  «        | j                   j                  «       ƒ d{  –—†  d| _         d| _        yy7 Œ# t        $ r Y Œw xY w­w©zClose the async connection.N)r   rK   Úwait_closedr9   r   rL   s    r   rX   zDirtyClient._close_async3  sf   è ø€ à�<‰<Ð#ðØ—‘×"Ñ"Ô$Ø—l‘l×.Ñ.Ó0×0Ð0ð  ˆDŒLØˆD�Lð $ð 1ùÜò Ùðüs:   ‚A-�7A ÁAÁA ÁA-ÁA Á	A*Á'A-Á)A*Á*A-c              ƒ   ó@   K  — | j                  «       ƒ d{  –—†  y7 Œ­wr^   )rX   rL   s    r   Úclose_asynczDirtyClient.close_async>  s   è ø€ à×ÑÓ!×!Ò!úó   ‚–—c                 ó&   — | j                  «        | S ©N)r   rL   s    r   Ú	__enter__zDirtyClient.__enter__F  s   € Ø�‰ŒØˆr   c                 ó$   — | j                  «        y rd   )rK   ©r   Úexc_typeÚexc_valÚexc_tbs       r   Ú__exit__zDirtyClient.__exit__J  s   € Ø�
‰
�r   c              ƒ   óB   K  — | j                  «       ƒ d {  –—†  | S 7 Œ­wrd   )rT   rL   s    r   Ú
__aenter__zDirtyClient.__aenter__M  s"   è ø€ Ø× Ñ Ó"×"Ð"Øˆð 	#ús   ‚–—c              ƒ   ó@   K  — | j                  «       ƒ d {  –—†  y 7 Œ­wrd   )ra   rg   s       r   Ú	__aexit__zDirtyClient.__aexit__Q  s   è ø€ Ø×ÑÓ × Ò úrb   N©ç      >@)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r+   r%   r@   r7   r8   rK   rT   rY   r\   rX   ra   re   rk   rm   ro   © r   r   r
   r
      sd   „ ñó&ò&ò*Hò*#KòJIò8Còò!òò46KòpNò8	 ò"òòòó!r   r
   c                   ó:   — e Zd ZdZdZdZ	 d
d„Zd„ Zd„ Zd„ Z	d	„ Z
y)r?   a  
    Iterator for streaming responses from dirty workers (sync).

    This class is returned by `DirtyClient.stream()` and yields chunks
    from a streaming response until the end message is received.

    Uses a deadline-based timeout approach:
    - Total stream timeout: limits entire stream duration
    - Idle timeout: limits gap between chunks (defaults to total timeout)
    rq   ç      @Nc                 óî   — || _         || _        || _        || _        || _        d| _        d| _        d | _        d | _        d | _	        |�|| _        y t        | j                  |j                  «      | _        y ©NF©Úclientr'   r(   r)   r*   Ú_startedÚ
_exhaustedÚ_request_idÚ	_deadlineÚ_last_chunk_timeÚminÚDEFAULT_IDLE_TIMEOUTr   Ú_idle_timeout©r   r|   r'   r(   r)   r*   Úidle_timeouts          r   r   zDirtyStreamIterator.__init__m  óx   € àˆŒØ ˆŒØˆŒØˆŒ	ØˆŒØˆŒØˆŒØˆÔØˆŒØ $ˆÔð )Ð4ˆLð 	Õä�T×.Ñ.°·±Ó?ð 	Õr   c                 ó   — | S rd   rv   rL   s    r   Ú__iter__zDirtyStreamIterator.__iter__  ó   € Øˆr   c                 óŒ   — | j                   rt        ‚| j                  s| j                  «        d| _        | j	                  «       S ©NT)r~   ÚStopIterationr}   Ú_start_requestÚ_read_next_chunkrL   s    r   Ú__next__zDirtyStreamIterator.__next__‚  s8   € Ø�?Š?ÜÐà�}Š}Ø×ÑÔ!Ø ˆDŒMà×$Ñ$Ó&Ð&r   c                 óH  — | j                   j                  5  | j                   j                  €| j                   j                  «        t	        j
                  «       }|| j                   j                  z   | _        || _        t        t        j                  «       «      | _        t        | j                  | j                  | j                  | j                   | j"                  ¬«      }t%        j&                  | j                   j                  |«       ddd«       y# 1 sw Y   yxY w©z(Send the initial request to the arbiter.N)r)   r*   )r|   r   r   r   ÚtimeÚ	monotonicr   r€   r�   r2   r3   r4   r   r   r'   r(   r)   r*   r   r5   ©r   Únowr<   s      r   rŽ   z"DirtyStreamIterator._start_requestŒ  sÒ   € à�[‰[×Ññ 	DØ�{‰{× Ñ Ð(Ø—‘×#Ñ#Ô%ô —.‘.Ó"ˆCØ  4§;¡;×#6Ñ#6Ñ6ˆDŒNØ$'ˆDÔ!ä"¤4§:¡:£<Ó0ˆDÔÜ"Ø× Ñ Ø—‘Ø—‘Ø—Y‘YØ—{‘{ôˆGô ×'Ñ'¨¯©×(9Ñ(9¸7ÔC÷#	D÷ 	Dñ 	Dús   —C8DÄD!c                 óx  — | j                   j                  5  t        j                  «       }|| j                  k\  r(d| _        t        d| j                   j                  ¬«      ‚| j                  |z
  }|| j                  kD  r| j                  }nt        || j                  «      }	 | j                   j                  j                  |«       t        j                  | j                   j                  «      }t        j                  «       | _        |j)                  d	«      }|t        j*                  k(  r|j)                  d
«      cddd«       S |t        j,                  k(  rd| _        t.        ‚|t        j0                  k(  r.d| _        |j)                  di «      }t3        j4                  |«      ‚|t        j6                  k(  rd| _        t.        ‚d| _        t3        d|› �«      ‚# t        j                  $ r~ t        j                  «       }|| j                  k\  r(d| _        t        d| j                   j                  ¬«      ‚|| j                   z
  }d| _        t        d|d›d�| j                  ¬«      ‚t"        $ r5}d| _        | j                   j%                  «        t'        d|› �«      |‚d}~ww xY w# 1 sw Y   yxY w)ú&Read the next message from the stream.TúStream exceeded total timeoutr0   ú%Timeout waiting for next chunk (idle ú.1fús)r1   NrB   Údatar    úUnknown message type: )r|   r   r“   r”   r€   r~   r   r   Ú_TIMEOUT_THRESHOLDr‚   r„   r   r   r   r6   r   r�   r9   r8   r   rD   ÚMSG_TYPE_CHUNKÚMSG_TYPE_ENDr�   rF   r   rG   rE   )	r   r–   Ú	remainingÚread_timeoutr=   Úidle_durationr#   rH   rI   s	            r   r�   z$DirtyStreamIterator._read_next_chunk¡  sy  € à�[‰[×Ññ F	Bä—.‘.Ó"ˆCØ�d—n‘nÒ$Ø"&�”Ü'Ø3Ø ŸK™K×/Ñ/ôð ð
 Ÿ™¨Ñ,ˆIð ˜4×2Ñ2Ò2Ø#×6Ñ6‘ä" 9¨d×.@Ñ.@ÓA�ðOØ—‘×!Ñ!×,Ñ,¨\Ô:Ü(×5Ñ5°d·k±k×6GÑ6GÓH�ô, %)§N¡NÓ$4ˆDÔ!à—|‘| FÓ+ˆHð œ=×7Ñ7Ò7Ø—|‘| FÓ+÷cF	Bñ F	Bðh œ=×5Ñ5Ò5Ø"&�”Ü#Ð#ð œ=×7Ñ7Ò7Ø"&�”Ø%Ÿ\™\¨'°2Ó6�
Ü ×*Ñ*¨:Ó6Ð6ð œ=×:Ñ:Ò:Ø"&�”ä#Ð#ð #ˆDŒOÜÐ5°h°ZÐ@ÓAÐAøôa —>‘>ò ä—n‘nÓ&�Ø˜$Ÿ.™.Ò(Ø&*�D”OÜ+Ø7Ø $§¡× 3Ñ 3ôð ð !$ d×&;Ñ&;Ñ ;�Ø"&�”Ü'Ø;¸MÈ#Ð;NÈbÐQØ ×.Ñ.ôð ô ò OØ"&�”Ø—‘×)Ñ)Ô+Ü*Ð-BÀ1À#Ð+FÓGÈQÐNûðOú÷KF	Bð F	Bús:   —BJ0Â%AG Ã3AJ0Å
BJ0Ç BJ-É80J(Ê(J-Ê-J0Ê0J9rd   )rr   rs   rt   ru   rƒ   rŸ   r   r‰   r�   rŽ   r�   rv   r   r   r?   r?   Z  s8   „ ñ	ð  Ðð Ðð #ó
ò$ò'òDó*HBr   r?   c                   ó:   — e Zd ZdZdZ	 d
d„Zd„ Zd„ Zd„ ZdZ	d	„ Z
y)r[   aÜ  
    Async iterator for streaming responses from dirty workers.

    This class is returned by `DirtyClient.stream_async()` and yields chunks
    from a streaming response until the end message is received.

    Uses a deadline-based timeout approach for efficiency:
    - Total stream timeout: limits entire stream duration
    - Idle timeout: limits gap between chunks (defaults to total timeout)

    This avoids the overhead of asyncio.wait_for() on every chunk read.
    rq   Nc                 óî   — || _         || _        || _        || _        || _        d| _        d| _        d | _        d | _        d | _	        |�|| _        y t        | j                  |j                  «      | _        y rz   r{   r…   s          r   r   z!DirtyAsyncStreamIterator.__init__ý  r‡   r   c                 ó   — | S rd   rv   rL   s    r   Ú	__aiter__z"DirtyAsyncStreamIterator.__aiter__  rŠ   r   c              ƒ   ó¼   K  — | j                   rt        ‚| j                  s| j                  «       ƒ d {  –—†  d| _        | j	                  «       ƒ d {  –—† S 7 Œ#7 Œ­wrŒ   )r~   ÚStopAsyncIterationr}   rŽ   r�   rL   s    r   Ú	__anext__z"DirtyAsyncStreamIterator.__anext__  sP   è ø€ Ø�?Š?Ü$Ð$à�}Š}Ø×%Ñ%Ó'×'Ð'Ø ˆDŒMà×*Ñ*Ó,×,Ð,ð (øð -ús!   ‚2A´AµAÁAÁAÁAc              ƒ   ó"  K  — | j                   j                  €"| j                   j                  «       ƒ d{  –—†  t        j                  «       }|| j                   j
                  z   | _        || _        t        t        j                  «       «      | _        t        | j                  | j                  | j                  | j                  | j                   ¬«      }t#        j$                  | j                   j                  |«      ƒ d{  –—†  y7 ŒÔ7 Œ­wr’   )r|   r   rT   r“   r”   r   r€   r�   r2   r3   r4   r   r   r'   r(   r)   r*   r   rV   r•   s      r   rŽ   z'DirtyAsyncStreamIterator._start_request  sÈ   è ø€ à�;‰;×ÑÐ&Ø—+‘+×+Ñ+Ó-×-Ð-ô �n‰nÓˆØ˜tŸ{™{×2Ñ2Ñ2ˆŒØ #ˆÔäœtŸz™z›|Ó,ˆÔÜØ×ÑØ�M‰MØ�K‰KØ—‘Ø—;‘;ô
ˆô ×/Ñ/°·±×0CÑ0CÀWÓM×MÑMð .øð 	Nús"   ‚4D¶D·CDÄDÄDÄDrx   c              ƒ   óp  K  — t        j                  «       }|| j                  k\  r(d| _        t	        d| j
                  j                  ¬«      ‚| j                  |z
  }	 || j                  kD  r2t        j                  | j
                  j                  «      ƒ d{  –—† }n\t        || j                  «      }t        j                  t        j                  | j
                  j                  «      |¬«      ƒ d{  –—† }t        j                  «       | _        |j)                  d	«      }|t        j*                  k(  r|j)                  d
«      S |t        j,                  k(  rd| _        t.        ‚|t        j0                  k(  r.d| _        |j)                  di «      }t3        j4                  |«      ‚|t        j6                  k(  rd| _        t.        ‚d| _        t3        d|› �«      ‚7 �ŒF7 Œë# t        j                  $ rw d| _        t        j                  «       }|| j                  k\  r!t	        d| j
                  j                  ¬«      ‚|| j                   z
  }t	        d|d›d�| j                  ¬«      ‚t"        $ r>}d| _        | j
                  j%                  «       ƒ d{  –—†7   t'        d|› �«      |‚d}~ww xY w­w)r˜   Tr™   r0   Nrš   r›   rœ   r1   rB   r�   r    rž   )r“   r”   r€   r~   r   r|   r   rŸ   r   rW   r   r‚   r„   rO   rP   rR   r�   r9   rX   r   rD   r    r¡   rª   rF   r   rG   rE   )	r   r–   r¢   r=   r£   r¤   r#   rH   rI   s	            r   r�   z)DirtyAsyncStreamIterator._read_next_chunk4  sr  è ø€ ô �n‰nÓˆð �$—.‘.Ò Ø"ˆDŒOÜ#Ø/ØŸ™×+Ñ+ôð ð
 —N‘N SÑ(ˆ	ð	Kð ˜4×2Ñ2Ò2Ü!.×!AÑ!AØ—K‘K×'Ñ'ó"÷ ‘ô
  # 9¨d×.@Ñ.@ÓA�Ü!(×!1Ñ!1Ü!×4Ñ4°T·[±[×5HÑ5HÓIØ(ô"÷ �ô. !%§¡Ó 0ˆÔà—<‘< Ó'ˆð ”}×3Ñ3Ò3Ø—<‘< Ó'Ð'ð ”}×1Ñ1Ò1Ø"ˆDŒOÜ$Ð$ð ”}×3Ñ3Ò3Ø"ˆDŒOØ!Ÿ™ g¨rÓ2ˆJÜ×&Ñ& zÓ2Ð2ð ”}×6Ñ6Ò6Ø"ˆDŒOÜ$Ð$ð ˆŒÜÐ1°(°Ð<Ó=Ð=ðoùðùô ×#Ñ#ò 	Ø"ˆDŒOÜ—.‘.Ó"ˆCØ�d—n‘nÒ$Ü'Ø3Ø ŸK™K×/Ñ/ôð ð   $×"7Ñ"7Ñ7ˆMÜ#Ø7¸ÀcÐ7JÈ"ÐMØ×*Ñ*ôð ô ò 	KØ"ˆDŒOØ—+‘+×*Ñ*Ó,×,Ñ,Ü&Ð)>¸q¸cÐ'BÓCÈÐJûð	Küs]   ‚AJ6Á;G$ ÂGÂAG$ Ã6G"Ã7G$ Ã;C$J6ÇG$ Ç"G$ Ç$BJ3É5$J.ÊJÊJ.Ê.J3Ê3J6rd   )rr   rs   rt   ru   rƒ   r   r¨   r«   rŽ   rŸ   r�   rv   r   r   r[   r[   ì  s7   „ ñð  Ðð #ó
ò$ò-òNð, ÐóJ>r   r[   Údirty_clientÚ_async_client_varc                 ó$   — | a ddlm}  || «       y)z@Set the global dirty socket path (called during initialization).r   )Úset_stash_socket_pathN)Ú_dirty_socket_pathÚstashr±   )Úpathr±   s     r   Úset_dirty_socket_pathrµ   ‘  s   € ð Ðõ -Ù˜$Õr   c                  óv   — t         €.t        j                  j                  d«      } | r| S t	        d«      ‚t         S )zGet the dirty socket path.ÚGUNICORN_DIRTY_SOCKETz\Dirty socket path not configured. Make sure dirty_workers > 0 and dirty_apps are configured.)r²   ÚosÚenvironrD   r   )r´   s    r   Úget_dirty_socket_pathrº   ›  s>   € äÐ!ä�z‰z�~‰~Ð5Ó6ˆÙØˆKÜðIó
ð 	
ô Ðr   Úreturnc                 óp   — t        t        dd«      }|€"t        «       }t        || ¬«      }|t        _        |S )aæ  
    Get or create a thread-local sync client.

    This is the recommended way to get a client in sync HTTP workers.

    Args:
        timeout: Timeout for operations in seconds

    Returns:
        DirtyClient: Thread-local client instance

    Example::

        from gunicorn.dirty import get_dirty_client

        def my_view(request):
            client = get_dirty_client()
            result = client.execute("myapp.ml:MLApp", "inference", data)
            return result
    r®   Nr0   )ÚgetattrÚ_thread_localrº   r
   r®   ©r   r|   r   s      r   Úget_dirty_clientrÀ   ©  s8   € ô* ”] N°DÓ9€FØ€~Ü+Ó-ˆÜ˜[°'Ô:ˆØ%+ŒÔ"Ø€Mr   c              ƒ   ó°   K  — 	 t         j                  «       }|S # t        $ r0 t        «       }t	        || ¬«      }t         j                  |«       Y |S w xY w­w)a  
    Get or create a context-local async client.

    This is the recommended way to get a client in async HTTP workers.

    Args:
        timeout: Timeout for operations in seconds

    Returns:
        DirtyClient: Context-local client instance

    Example::

        from gunicorn.dirty import get_dirty_client_async

        async def my_view(request):
            client = await get_dirty_client_async()
            result = await client.execute_async("myapp.ml:MLApp", "inference", data)
            return result
    r0   )r¯   rD   ÚLookupErrorrº   r
   Úsetr¿   s      r   Úget_dirty_client_asyncrÄ   Æ  sW   è ø€ ð*&Ü"×&Ñ&Ó(ˆð
 €Møô	 ò &Ü+Ó-ˆÜ˜[°'Ô:ˆÜ×Ñ˜fÕ%Ø€Mð	&üs%   ‚A„ ˜Aš5AÁAÁAÁAc                  ób   — t        t        dd«      } | �| j                  «        dt        _        yy)z4Close the thread-local client (call on worker exit).r®   N)r½   r¾   rK   r®   ©r|   s    r   Úclose_dirty_clientrÇ   ä  s,   € ä”] N°DÓ9€FØÐØ�‰ŒØ%)ŒÕ"ð r   c               ƒ   óˆ   K  — 	 t         j                  «       } | j                  «       ƒ d{  –—†  y7 Œ# t        $ r Y yw xY w­w)z%Close the context-local async client.N)r¯   rD   ra   rÂ   rÆ   s    r   Úclose_dirty_client_asyncrÉ   ì  s<   è ø€ ðÜ"×&Ñ&Ó(ˆØ× Ñ Ó"×"Ò"ùÜò Ùðüs,   ‚A„'3 «1¬3 °A±3 ³	?¼A¾?¿Arp   )ru   rO   Úcontextvarsr¸   r   r   r“   r3   Úerrorsr   r   r   Úprotocolr   r   r
   r?   r[   Úlocalr¾   Ú
ContextVarr¯   Ú__annotations__r²   rµ   rº   rÀ   rÄ   rÇ   rÉ   rv   r   r   ú<module>rÐ      sÈ   ðò
ó Û Û 	Û Û Û Û ÷ñ ÷
÷s!ñ s!÷v	OBñ OB÷dR>ñ R>ðt  �	—‘Ó!€ð :P¸×9OÑ9OØó:Ð �;×)Ñ)¨+Ñ6ó ð
 Ð ò òñ kó ñ:°+ó ò<*ór   