§
    …ßjì^  ã                   óL  — 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dS )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                   óz   — 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dS )Ú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        ¦   «         | _        dS )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      úH/var/www/html/venv/lib/python3.11/site-packages/gunicorn/dirty/client.pyÚ__init__zDirtyClient.__init__(   s;   € ð 'ˆÔØˆŒØˆŒ
ØˆŒØˆŒÝ”^Ñ%Ô%ˆŒ
ˆ
ˆ
ó    c                 óp  — | j         �dS 	 t          j        t          j        t          j        ¦  «        | _         | j                              | j        ¦  «         | j                              | j        ¦  «         dS # 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¿   € ð Œ:Ð!ØˆFð		Ýœ¥v¤~µvÔ7IÑJÔJˆDŒJØŒJ×!Ò! $¤,Ñ/Ô/Ð/ØŒJ×Ò˜tÔ/Ñ0Ô0Ð0Ð0Ð0øÝ”�gÐ&ð 	ð 	ð 	ØˆDŒJÝ&Ø:°qÐ:Ð:Ø Ô,ðñ ô ð ðøøøøð	øøøs   ‹A,A9 Á9B5Â!B0Â0B5c                 ót   — | j         5  |                      ||||¦  «        cddd¦  «         S # 1 swxY w Y   dS )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   s›   € ð$ ŒZð 	Hð 	HØ×'Ò'¨°&¸$ÀÑGÔGð	Hð 	Hð 	Hð 	Hñ 	Hô 	Hð 	Hð 	Hð 	Hð 	Hð 	Hð 	Høøøð 	Hð 	Hð 	Hð 	Hð 	Hð 	Hs   ˆ-­1´1c                 ó@  — | j         €|                      ¦   «          t          t          j        ¦   «         ¦  «        }t          |||||¬¦  «        }	 t          j        | j         |¦  «         t          j        | j         ¦  «        }|  	                    |¦  «        S # t          j        $ r+ |                      ¦   «          t          d| j        ¬¦  «        ‚t          $ rB}|                      ¦   «          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ŒNˆNõ �œ™œÑ&Ô&ˆ
ÝØ!ØØØØð
ñ 
ô 
ˆð	KåÔ'¨¬
°GÑ<Ô<Ð<õ %Ô1°$´*Ñ=Ô=ˆHð ×(Ò(¨Ñ2Ô2Ð2øÝŒ~ð 	ð 	ð 	Ø×ÒÑ Ô Ð Ý#Ø8Øœðñ ô ð õ ð 	Kð 	Kð 	KØ×ÒÑ Ô Ð Ý˜!�ZÑ(Ô(ð ØÝ&Ð'B¸qÐ'BÐ'BÑCÔCÈÐJøøøøð		Køøøs   ÁAB ÂADÃ=DÄ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ÔHÐHr   c                 ó  — |                      d¦  «        }|t          j        k    r|                      d¦  «        S |t          j        k    r,|                      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   r8   zDirtyClient._handle_response¦   s€   € à—<’< Ñ'Ô'ˆà•}Ô6Ò6Ð6Ø—<’< Ñ)Ô)Ð)Ø�Ô5Ò5Ð5Ø!Ÿš g¨rÑ2Ô2ˆJÝÔ(¨Ñ4Ô4ˆEØˆKåÐA°xÐAÐAÑBÔBÐBr   c                 ó|   — | j         �4	 | j                              ¦   «          n# t          $ r Y nw xY wd| _         dS dS )zClose the socket connection.N)r   Úcloser:   ©r   s    r   r9   zDirtyClient._close_socket³   sY   € àŒ:Ð!ðØ”
× Ò Ñ"Ô"Ð"Ð"øÝð ð ð Ø�ðøøøàˆDŒJˆJˆJð "Ð!s   ‰# £
0¯0c                 ón   — | j         5  |                      ¦   «          ddd¦  «         dS # 1 swxY w Y   dS )zClose the sync connection.N)r   r9   rM   s    r   rL   zDirtyClient.close¼   s€   € àŒZð 	!ð 	!Ø×ÒÑ Ô Ð ð	!ð 	!ð 	!ñ 	!ô 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!øøøð 	!ð 	!ð 	!ð 	!ð 	!ð 	!s   ˆ*ª.±.c              ƒ   óV  K  — | j         �dS 	 t          j        t          j        | j        ¦  «        | j        ¬¦  «        ƒ d{V —†\  | _        | _         dS # t          j        $ r t          d| j        ¬¦  «        ‚t          t          f$ r}t          d|› �| j        ¬¦  «        |‚d}~ww xY w)z
        Establish async connection to arbiter.

        Raises:
            DirtyConnectionError: If connection fails
        Nr1   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û   è è € ð Œ<Ð#ØˆFð	Ý/6Ô/?ÝÔ,¨TÔ-=Ñ>Ô>Øœð0ñ 0ô 0ð *ð *ð *ð *ð *ð *Ñ&ˆDŒL˜$œ,˜,˜,øõ Ô#ð 	ð 	ð 	Ý#Ø5Øœðñ ô ð õ �Ð)ð 	ð 	ð 	Ý&Ø:°qÐ:Ð:Ø Ô,ðñ ô ð ðøøøøð	øøøs   �AA Á5B(Â	B#Â#B(c              �   ó²  K  — | j         €|                      ¦   «         ƒ d{V —† t          t          j        ¦   «         ¦  «        }t          |||||¬¦  «        }	 t          j        | j         |¦  «        ƒ d{V —† t          j	        t          j
        | j        ¦  «        | j        ¬¦  «        ƒ d{V —†}|                      |¦  «        S # t          j        $ r1 |                      ¦   «         ƒ d{V —† t!          d| j        ¬¦  «        ‚t"          $ rH}|                      ¦   «         ƒ d{V —† t%          |t&          ¦  «        r‚ t)          d|› �¦  «        |‚d}~ww xY 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.   r1   r0   r2   )r   rU   r3   r4   r5   r   r   Úwrite_message_asyncrP   rQ   Úread_message_asyncr   r   r8   rS   Ú_close_asyncr   r:   r;   r   r   r<   s	            r   Úexecute_asynczDirtyClient.execute_asyncß   sÊ  è è € ð& Œ<ÐØ×$Ò$Ñ&Ô&Ð&Ð&Ð&Ð&Ð&Ð&Ð&õ �œ™œÑ&Ô&ˆ
ÝØ!ØØØØð
ñ 
ô 
ˆð	KåÔ3°D´LÀ'ÑJÔJÐJÐJÐJÐJÐJÐJÐJõ %Ô-ÝÔ0°´Ñ>Ô>Øœðñ ô ð ð ð ð ð ð ˆHð ×(Ò(¨Ñ2Ô2Ð2øÝÔ#ð 	ð 	ð 	Ø×#Ò#Ñ%Ô%Ð%Ð%Ð%Ð%Ð%Ð%Ð%Ý#Ø8Øœðñ ô ð õ ð 	Kð 	Kð 	KØ×#Ò#Ñ%Ô%Ð%Ð%Ð%Ð%Ð%Ð%Ð%Ý˜!�ZÑ(Ô(ð ØÝ&Ð'B¸qÐ'BÐ'BÑCÔCÈÐJøøøøð		Køøøs   ÁA,C ÃAEÄA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ÔMÐMr   c              ƒ   óÌ   K  — | j         �Z	 | j                              ¦   «          | j                              ¦   «         ƒ d{V —† n# t          $ r Y nw xY wd| _         d| _        dS dS ©zClose the async connection.N)r   rL   Úwait_closedr:   r   rM   s    r   rY   zDirtyClient._close_async3  sŠ   è è € àŒ<Ð#ðØ”×"Ò"Ñ$Ô$Ð$Ø”l×.Ò.Ñ0Ô0Ð0Ð0Ð0Ð0Ð0Ð0Ð0Ð0øÝð ð ð Ø�ðøøøàˆDŒLØˆDŒLˆLˆLð $Ð#s   ‹8A Á
AÁAc              ƒ   ó>   K  — |                       ¦   «         ƒ d{V —† dS r_   )rY   rM   s    r   Úclose_asynczDirtyClient.close_async>  s0   è è € à×ÒÑ!Ô!Ð!Ð!Ð!Ð!Ð!Ð!Ð!Ð!Ð!r   c                 ó.   — |                       ¦   «          | S ©N)r    rM   s    r   Ú	__enter__zDirtyClient.__enter__F  s   € Ø�Š‰ŒˆØˆr   c                 ó.   — |                       ¦   «          d S rd   )rL   ©r   Úexc_typeÚexc_valÚexc_tbs       r   Ú__exit__zDirtyClient.__exit__J  s   € Ø�
Š
‰Œˆˆˆr   c              ƒ   ó>   K  — |                       ¦   «         ƒ d {V —† | S rd   )rU   rM   s    r   Ú
__aenter__zDirtyClient.__aenter__M  s/   è è € Ø× Ò Ñ"Ô"Ð"Ð"Ð"Ð"Ð"Ð"Ð"Øˆr   c              ƒ   ó>   K  — |                       ¦   «         ƒ d {V —† d S rd   )rb   rg   s       r   Ú	__aexit__zDirtyClient.__aexit__Q  s0   è è € Ø×ÒÑ Ô Ð Ð Ð Ð Ð Ð Ð Ð Ð r   N©r   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r    r,   r&   rA   r8   r9   rL   rU   rZ   r]   rY   rb   re   rk   rm   ro   © r   r   r
   r
      s6  € € € € € ðð ð&ð &ð &ð &ð&ð ð ð*Hð Hð Hð*#Kð #Kð #KðJIð Ið Ið8Cð Cð Cðð ð ð!ð !ð !ðð ð ð46Kð 6Kð 6KðpNð Nð 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
dS )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)
    r   ç      @Nc                 óØ   — || _         || _        || _        || _        || _        d| _        d| _        d | _        d | _        d | _	        |�|nt          | j        |j        ¦  «        | _        d S ©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  óv   € àˆŒØ ˆŒØˆŒØˆŒ	ØˆŒØˆŒØˆŒØˆÔØˆŒØ $ˆÔð )Ð4ˆLˆLÝ�TÔ.°´Ñ?Ô?ð 	ÔÐÐr   c                 ó   — | S rd   ru   rM   s    r   Ú__iter__zDirtyStreamIterator.__iter__  ó   € Øˆr   c                 óŠ   — | j         rt          ‚| j        s|                      ¦   «          d| _        |                      ¦   «         S ©NT)r}   ÚStopIterationr|   Ú_start_requestÚ_read_next_chunkrM   s    r   Ú__next__zDirtyStreamIterator.__next__‚  sG   € ØŒ?ð 	 ÝÐàŒ}ð 	!Ø×ÒÑ!Ô!Ð!Ø ˆDŒMà×$Ò$Ñ&Ô&Ð&r   c                 óØ  — | j         j        5  | j         j        €| j                              ¦   «          t	          j        ¦   «         }|| j         j        z   | _        || _        t          t          j        ¦   «         ¦  «        | _        t          | j        | j        | j        | j        | j        ¬¦  «        }t%          j        | j         j        |¦  «         ddd¦  «         dS # 1 swxY w Y   dS ©z(Send the initial request to the arbiter.N)r*   r+   )r{   r   r   r    ÚtimeÚ	monotonicr   r   r€   r3   r4   r5   r~   r   r(   r)   r*   r+   r   r6   ©r   Únowr=   s      r   r�   z"DirtyStreamIterator._start_requestŒ  s*  € àŒ[Ôð 	Dð 	DØŒ{Ô Ð(Ø”×#Ò#Ñ%Ô%Ð%õ ”.Ñ"Ô"ˆCØ  4¤;Ô#6Ñ6ˆDŒNØ$'ˆDÔ!å"¥4¤:¡<¤<Ñ0Ô0ˆDÔÝ"ØÔ Ø”Ø”Ø”YØ”{ðñ ô ˆGõ Ô'¨¬Ô(9¸7ÑCÔCÐCð#	Dð 	Dð 	Dñ 	Dô 	Dð 	Dð 	Dð 	Dð 	Dð 	Dð 	Dð 	Døøøð 	Dð 	Dð 	Dð 	Dð 	Dð 	Ds   �CCÃC#Ã&C#c                 óþ  — | j         j        5  t          j        ¦   «         }|| j        k    r"d| _        t          d| j         j        ¬¦  «        ‚| j        |z
  }|| j        k    r| j        }nt          || j
        ¦  «        }	 | j         j                             |¦  «         t          j        | j         j        ¦  «        }n¿# t          j        $ rm t          j        ¦   «         }|| j        k    r"d| _        t          d| j         j        ¬¦  «        ‚|| j        z
  }d| _        t          d|d›d�| j
        ¬¦  «        ‚t"          $ r8}d| _        | j                              ¦   «          t'          d|› �¦  «        |‚d}~ww xY wt          j        ¦   «         | _        |                     d	¦  «        }|t          j        k    r!|                     d
¦  «        cddd¦  «         S |t          j        k    rd| _        t.          ‚|t          j        k    r1d| _        |                     di ¦  «        }t3          j        |¦  «        ‚|t          j        k    rd| _        t.          ‚d| _        t3          d|› �¦  «        ‚# 1 swxY w Y   dS )ú&Read the next message from the stream.TúStream exceeded total timeoutr1   ú%Timeout waiting for next chunk (idle ú.1fús)r2   NrC   Údatar!   úUnknown message type: )r{   r   r’   r“   r   r}   r   r   Ú_TIMEOUT_THRESHOLDr�   rƒ   r   r   r   r7   r   r€   r:   r9   r   rE   ÚMSG_TYPE_CHUNKÚMSG_TYPE_ENDrŒ   rG   r   rH   rF   )	r   r•   Ú	remainingÚread_timeoutr>   Úidle_durationr$   rI   rJ   s	            r   rŽ   z$DirtyStreamIterator._read_next_chunk¡  s  € àŒ[Ôð F	Bð F	Bå”.Ñ"Ô"ˆCØ�d”nÒ$Ð$Ø"&�”Ý'Ø3Ø œKÔ/ðñ ô ð ð
 œ¨Ñ,ˆIð ˜4Ô2Ò2Ð2Ø#Ô6��å" 9¨dÔ.@ÑAÔA�ðOØ”Ô!×,Ò,¨\Ñ:Ô:Ð:Ý(Ô5°d´kÔ6GÑHÔH��øÝ”>ð ð ð å”nÑ&Ô&�Ø˜$œ.Ò(Ð(Ø&*�D”OÝ+Ø7Ø $¤Ô 3ðñ ô ð ð !$ dÔ&;Ñ ;�Ø"&�”Ý'ØQ¸MÐQÐQÐQÐQØ Ô.ðñ ô ð õ ð Oð Oð OØ"&�”Ø”×)Ò)Ñ+Ô+Ð+Ý*Ð+FÀ1Ð+FÐ+FÑGÔGÈQÐNøøøøðOøøøõ %)¤NÑ$4Ô$4ˆDÔ!à—|’| FÑ+Ô+ˆHð �=Ô7Ò7Ð7Ø—|’| FÑ+Ô+ðcF	Bð F	Bð F	Bð F	Bñ F	Bô F	Bð F	Bð F	Bðh �=Ô5Ò5Ð5Ø"&�”Ý#Ð#ð �=Ô7Ò7Ð7Ø"&�”Ø%Ÿ\š\¨'°2Ñ6Ô6�
Ý Ô*¨:Ñ6Ô6Ð6ð �=Ô:Ò:Ð:Ø"&�”å#Ð#ð #ˆDŒOÝÐ@°hÐ@Ð@ÑAÔAÐAðMF	Bð F	Bð F	Bð F	Bøøøð F	Bð F	Bð F	Bð F	Bð F	Bð F	Bs?   �A3I2Â=B?Â>I2Â?BE;Å3E6Å6E;Å;AI2ÇBI2É2I6É9I6rd   )rq   rr   rs   rt   r‚   rž   r   rˆ   r�   r�   rŽ   ru   r   r   r@   r@   Z  s�   € € € € € ð	ð 	ð  Ðð Ðð #ð
ð 
ð 
ð 
ð$ð ð ð'ð 'ð 'ðDð Dð Dð*HBð HBð HBð HBð 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
dS )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.
    r   Nc                 óØ   — || _         || _        || _        || _        || _        d| _        d| _        d | _        d | _        d | _	        |�|nt          | j        |j        ¦  «        | _        d S ry   rz   r„   s          r   r   z!DirtyAsyncStreamIterator.__init__ý  r†   r   c                 ó   — | S rd   ru   rM   s    r   Ú	__aiter__z"DirtyAsyncStreamIterator.__aiter__  r‰   r   c              ƒ   ó¦   K  — | j         rt          ‚| j        s!|                      ¦   «         ƒ d {V —† d| _        |                      ¦   «         ƒ d {V —†S r‹   )r}   ÚStopAsyncIterationr|   r�   rŽ   rM   s    r   Ú	__anext__z"DirtyAsyncStreamIterator.__anext__  so   è è € ØŒ?ð 	%Ý$Ð$àŒ}ð 	!Ø×%Ò%Ñ'Ô'Ð'Ð'Ð'Ð'Ð'Ð'Ð'Ø ˆDŒMà×*Ò*Ñ,Ô,Ð,Ð,Ð,Ð,Ð,Ð,Ð,r   c              ƒ   óª  K  — | j         j        €| j                              ¦   «         ƒ d{V —† t          j        ¦   «         }|| j         j        z   | _        || _        t          t          j
        ¦   «         ¦  «        | _        t          | j        | j        | j        | j        | j        ¬¦  «        }t#          j        | j         j        |¦  «        ƒ d{V —† dS r‘   )r{   r   rU   r’   r“   r   r   r€   r3   r4   r5   r~   r   r(   r)   r*   r+   r   rW   r”   s      r   r�   z'DirtyAsyncStreamIterator._start_request  sÕ   è è € àŒ;ÔÐ&Ø”+×+Ò+Ñ-Ô-Ð-Ð-Ð-Ð-Ð-Ð-Ð-õ ŒnÑÔˆØ˜tœ{Ô2Ñ2ˆŒØ #ˆÔå�tœz™|œ|Ñ,Ô,ˆÔÝØÔØŒMØŒKØ”Ø”;ð
ñ 
ô 
ˆõ Ô/°´Ô0CÀWÑMÔMÐMÐMÐMÐMÐMÐMÐMÐMÐMr   rw   c              ƒ   óä  K  — t          j        ¦   «         }|| j        k    r"d| _        t	          d| j        j        ¬¦  «        ‚| j        |z
  }	 || j        k    r%t          j	        | j        j
        ¦  «        ƒ d{V —†}nMt          || j        ¦  «        }t          j        t          j	        | j        j
        ¦  «        |¬¦  «        ƒ d{V —†}n¾# t          j        $ rf d| _        t          j        ¦   «         }|| j        k    rt	          d| j        j        ¬¦  «        ‚|| j        z
  }t	          d|d›d�| j        ¬¦  «        ‚t"          $ r>}d| _        | j                             ¦   «         ƒ d{V —† t'          d|› �¦  «        |‚d}~ww xY wt          j        ¦   «         | _        |                     d	¦  «        }|t          j        k    r|                     d
¦  «        S |t          j        k    rd| _        t.          ‚|t          j        k    r1d| _        |                     di ¦  «        }t3          j        |¦  «        ‚|t          j        k    rd| _        t.          ‚d| _        t3          d|› �¦  «        ‚)r—   Tr˜   r1   Nr™   rš   r›   r2   rC   rœ   r!   r�   )r’   r“   r   r}   r   r{   r   rž   r   rX   r   r�   rƒ   rP   rQ   rS   r€   r:   rY   r   rE   rŸ   r    r©   rG   r   rH   rF   )	r   r•   r¡   r>   r¢   r£   r$   rI   rJ   s	            r   rŽ   z)DirtyAsyncStreamIterator._read_next_chunk4  sç  è è € õ ŒnÑÔˆð �$”.Ò Ð Ø"ˆDŒOÝ#Ø/ØœÔ+ðñ ô ð ð
 ”N SÑ(ˆ	ð	Kð ˜4Ô2Ò2Ð2Ý!.Ô!AØ”KÔ'ñ"ô "ð ð ð ð ð ð ��õ
  # 9¨dÔ.@ÑAÔA�Ý!(Ô!1Ý!Ô4°T´[Ô5HÑIÔIØ(ð"ñ "ô "ð ð ð ð ð ð �øøõ Ô#ð 	ð 	ð 	Ø"ˆDŒOÝ”.Ñ"Ô"ˆCØ�d”nÒ$Ð$Ý'Ø3Ø œKÔ/ðñ ô ð ð   $Ô"7Ñ7ˆMÝ#ØM¸ÐMÐMÐMÐMØÔ*ðñ ô ð õ ð 	Kð 	Kð 	KØ"ˆDŒOØ”+×*Ò*Ñ,Ô,Ð,Ð,Ð,Ð,Ð,Ð,Ð,Ý&Ð'B¸qÐ'BÐ'BÑCÔCÈÐJøøøøð	Køøøõ !%¤Ñ 0Ô 0ˆÔà—<’< Ñ'Ô'ˆð •}Ô3Ò3Ð3Ø—<’< Ñ'Ô'Ð'ð •}Ô1Ò1Ð1Ø"ˆDŒOÝ$Ð$ð •}Ô3Ò3Ð3Ø"ˆDŒOØ!Ÿš g¨rÑ2Ô2ˆJÝÔ& zÑ2Ô2Ð2ð •}Ô6Ò6Ð6Ø"ˆDŒOÝ$Ð$ð ˆŒÝÐ<°(Ð<Ð<Ñ=Ô=Ð=s   ÁA=C ÃA=FÅ	9FÆFrd   )rq   rr   rs   rt   r‚   r   r§   rª   r�   rž   rŽ   ru   r   r   r\   r\   ì  s‹   € € € € € ðð ð  Ðð #ð
ð 
ð 
ð 
ð$ð ð ð-ð -ð -ðNð Nð Nð, ÐðJ>ð J>ð J>ð J>ð J>r   r\   Údirty_clientÚ_async_client_varc                 ó,   — | a ddlm}  || ¦  «         dS )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´   ‘  s5   € ð Ðð -Ð,Ð,Ð,Ð,Ð,ØÐ˜$ÑÔÐÐÐr   c                  ó‚   — t           €2t          j                             d¦  «        } | r| S t	          d¦  «        ‚t           S )zGet the dirty socket path.NÚGUNICORN_DIRTY_SOCKETz\Dirty socket path not configured. Make sure dirty_workers > 0 and dirty_apps are configured.)r±   ÚosÚenvironrE   r   )r³   s    r   Úget_dirty_socket_pathr¹   ›  sI   € åÐ!åŒz�~Š~Ð5Ñ6Ô6ˆØð 	ØˆKÝðIñ
ô 
ð 	
õ Ðr   r   Úreturnc                 óŒ   — 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­   Nr1   )ÚgetattrÚ_thread_localr¹   r
   r­   ©r   r{   r   s      r   Úget_dirty_clientr¿   ©  sD   € õ* •] N°DÑ9Ô9€FØ€~Ý+Ñ-Ô-ˆÝ˜[°'Ð:Ñ:Ô:ˆØ%+�Ô"Ø€Mr   c              ƒ   óÒ   K  — 	 t                                ¦   «         }nI# t          $ r< t          ¦   «         }t	          || ¬¦  «        }t                                |¦  «         Y nw xY w|S )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
    r1   )r®   rE   ÚLookupErrorr¹   r
   Úsetr¾   s      r   Úget_dirty_client_asyncrÃ   Æ  sw   è è € ð*&Ý"×&Ò&Ñ(Ô(ˆˆøÝð &ð &ð &Ý+Ñ-Ô-ˆÝ˜[°'Ð:Ñ:Ô:ˆÝ×Ò˜fÑ%Ô%Ð%Ð%Ð%ð&øøøð €Ms   „ žAA$Á#A$c                  óz   — t          t          dd¦  «        } | �"|                      ¦   «          dt          _        dS dS )z4Close the thread-local client (call on worker exit).r­   N)r¼   r½   rL   r­   ©r{   s    r   Úclose_dirty_clientrÆ   ä  s<   € å•] N°DÑ9Ô9€FØÐØ�Š‰ŒˆØ%)�Ô"Ð"Ð"ð Ðr   c               ƒ   ó”   K  — 	 t                                ¦   «         } |                      ¦   «         ƒ d{V —† dS # t          $ r Y dS w xY w)z%Close the context-local async client.N)r®   rE   rb   rÁ   rÅ   s    r   Úclose_dirty_client_asyncrÈ   ì  sh   è è € ðÝ"×&Ò&Ñ(Ô(ˆØ× Ò Ñ"Ô"Ð"Ð"Ð"Ð"Ð"Ð"Ð"Ð"Ð"øÝð ð ð Øˆˆðøøøs   „39 ¹
AÁArp   )rt   rP   Úcontextvarsr·   r   r   r’   r4   Úerrorsr   r   r   Úprotocolr   r   r
   r@   r\   Úlocalr½   Ú
ContextVarr®   Ú__annotations__r±   r´   r¹   r¿   rÃ   rÆ   rÈ   ru   r   r   ú<module>rÏ      s#  ðð
ð ð ð €€€Ø Ð Ð Ð Ø 	€	€	€	Ø €€€Ø Ð Ð Ð Ø €€€Ø €€€ðð ð ð ð ð ð ð ð ð ð
ð ð ð ð ð ð ð ðs!ð s!ð s!ð s!ð s!ñ s!ô s!ð s!ðv	OBð OBð OBð OBð OBñ OBô OBð OBðdR>ð R>ð R>ð R>ð R>ñ R>ô R>ð R>ðt  �	”Ñ!Ô!€ð :P¸Ô9OØñ:ô :Ð �;Ô)¨+Ô6ð ð ñ ð
 Ð ð ð  ð  ðð ð ðð  kð ð ð ð ð:ð °+ð ð ð ð ð<*ð *ð *ðð ð ð ð r   