
    \ei&                         d Z ddlZddlZddlZddlZddlmZ  G d d      Z	 ddeded	ed
e	de
de
fdZdede
fdZdede
fdZdede
fdZdede
fdZy)ah  
Dirty Arbiters Protocol

Length-prefixed JSON message framing over Unix sockets.
Provides both async (primary) and sync (for HTTP workers) APIs.

Message Format:
+----------------+------------------+
| 4-byte length  | JSON payload     |
+----------------+------------------+

The length field is a 4-byte unsigned integer in network byte order (big-endian).
    N   )DirtyProtocolErrorc                   v   e Zd ZdZdZ ej                  e      ZdZdZ	dZ
dZdZdZed	ed
efd       Zeded
efd       Zedej(                  d
efd       Zedej,                  d	ed
dfd       Zedej0                  ded
efd       Zedej0                  d
efd       Zedej0                  d	ed
dfd       Zy)DirtyProtocolz0Length-prefixed JSON messages over Unix sockets.z!Ii   requestresponseerrorchunkendmessagereturnc                    	 t        j                  |       j                  d      }t        |      t        j
                  kD  r)t        dt        |       dt        j
                   d      t        j                  t        j                  t        |            }||z   S # t        t        f$ r}t        d|       d}~ww xY w)a
  
        Encode a message dict to length-prefixed bytes.

        Args:
            message: Dictionary to encode as JSON

        Returns:
            bytes: Length-prefixed encoded message

        Raises:
            DirtyProtocolError: If encoding fails
        utf-8Message too large:  bytes (max: )zFailed to encode message: N)jsondumpsencodelenr   MAX_MESSAGE_SIZEr   structpackHEADER_FORMAT	TypeError
ValueError)r   payloadlengthes       M/var/www/ceduoda/envi/lib/python3.12/site-packages/gunicorn/dirty/protocol.pyr   zDirtyProtocol.encode,   s    
	Gjj)009G7|m<<<()#g, 8*;;<A?  [[!<!<c'lKFG##:& 	G$'A!%EFF	Gs   BB B?,B::B?datac                     	 t        j                  | j                  d            S # t         j                  t        f$ r}t        d| |       d}~ww xY w)z
        Decode bytes (without length prefix) to message dict.

        Args:
            data: JSON bytes to decode

        Returns:
            dict: Decoded message

        Raises:
            DirtyProtocolError: If decoding fails
        r   zFailed to decode message: raw_dataN)r   loadsdecodeJSONDecodeErrorUnicodeDecodeErrorr   )r!   r   s     r    r&   zDirtyProtocol.decodeF   sU    	4::dkk'233$$&89 	4$'A!%E.24 4	4s   #& AAAreaderc                   K   	 | j                  t        j                         d{   }t        j                  t        j                  |      d   }|t        j                  kD  r t        d| dt        j                   d      |dk(  rt        d	      	 | j                  |       d{   }t        j                  |      S 7 # t        j                  $ r\}t        |j                        dk(  r t        dt        |j                         dt        j                   |j                        d}~ww xY w7 # t        j                  $ r5}t        d
t        |j                         d| |j                        d}~ww xY ww)aF  
        Read a complete message from async stream.

        Args:
            reader: asyncio StreamReader

        Returns:
            dict: Decoded message

        Raises:
            DirtyProtocolError: If read fails or message is malformed
            asyncio.IncompleteReadError: If connection closed mid-read
        Nr   zIncomplete header: got  bytes, expected r#   r   r   r   Empty message receivedzIncomplete message: got )readexactlyr   HEADER_SIZEasyncioIncompleteReadErrorr   partialr   r   unpackr   r   r&   )r)   headerr   r   r   s        r    read_message_asyncz DirtyProtocol.read_message_async^   sr     
	!--m.G.GHHF }::FCAFM222$%fX .&778; 
 Q;$%=>>	"..v66G ##G,,A I** 	199~"$)#aii.)9 :)5568 		. 7** 	$*3qyy>*: ;"8% 	so   F"C CC A*FD: *D8+D: /FC D5AD00D55F8D: :F0E==FFwriterNc                    K   t         j                  |      }| j                  |       | j                          d{    y7 w)a  
        Write a message to async stream.

        Args:
            writer: asyncio StreamWriter
            message: Dictionary to send

        Raises:
            DirtyProtocolError: If encoding fails
            ConnectionError: If write fails
        N)r   r   writedrain)r5   r   r!   s      r    write_message_asyncz!DirtyProtocol.write_message_async   s3      ##G,Tllns   :AAAsocknc                     d}t        |      |k  rh| j                  |t        |      z
        }|s5t        |      dk(  rt        d      t        dt        |       d| |      ||z  }t        |      |k  rh|S )a  
        Receive exactly n bytes from a socket.

        Args:
            sock: Socket to read from
            n: Number of bytes to read

        Returns:
            bytes: Received data

        Raises:
            DirtyProtocolError: If read fails or connection closed
            r   zConnection closedzConnection closed after r+   r#   )r   recvr   )r:   r;   r!   r
   s       r    _recv_exactlyzDirtyProtocol._recv_exactly   s     $i!mIIa#d)m,Et9>,-@AA(.s4yk9J1#N!  EMD $i!m r=   c                 t   t         j                  | t         j                        }t        j                  t         j
                  |      d   }|t         j                  kD  r t        d| dt         j                   d      |dk(  rt        d      t         j                  | |      }t         j                  |      S )z
        Read a complete message from socket (sync).

        Args:
            sock: Socket to read from

        Returns:
            dict: Decoded message

        Raises:
            DirtyProtocolError: If read fails or message is malformed
        r   r   r   r   r,   )	r   r?   r.   r   r2   r   r   r   r&   )r:   r3   r   r   s       r    read_messagezDirtyProtocol.read_message   s     ,,T=3L3LM}::FCAFM222$%fX .&778; 
 Q;$%=>>  --dF;##G,,r=   c                 P    t         j                  |      }| j                  |       y)z
        Write a message to socket (sync).

        Args:
            sock: Socket to write to
            message: Dictionary to send

        Raises:
            DirtyProtocolError: If encoding fails
            OSError: If write fails
        N)r   r   sendall)r:   r   r!   s      r    write_messagezDirtyProtocol.write_message   s      ##G,Tr=   )__name__
__module____qualname____doc__r   r   calcsizer.   r   MSG_TYPE_REQUESTMSG_TYPE_RESPONSEMSG_TYPE_ERRORMSG_TYPE_CHUNKMSG_TYPE_ENDstaticmethoddictbytesr   r&   r/   StreamReaderr4   StreamWriterr9   socketintr?   rA   rD    r=   r    r   r      sa   : M!&//-0K ( !"NNLG G G G2 4U 4t 4 4. 0-)=)= 0-$ 0- 0-d '*>*> +/48 * FMM c e  6 -6== -T - -< FMM D T  r=   r   
request_idapp_pathactionargskwargsr   c                 R    t         j                  | |||rt        |      ng |xs i dS )a>  
    Build a request message.

    Args:
        request_id: Unique request identifier
        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:
        dict: Request message
    )typeidrX   rY   rZ   r[   )r   rJ   list)rW   rX   rY   rZ   r[   s        r    make_requestr`      s2      .."T
,B r=   c                 *    t         j                  | |dS )z
    Build a success response message.

    Args:
        request_id: Request identifier this responds to
        result: Result value (must be JSON-serializable)

    Returns:
        dict: Response message
    )r]   r^   result)r   rK   )rW   rb   s     r    make_responserc     s     // r=   c                     ddl m} t        ||      r|j                         }n5t        |t              r|}n"t        |      j                  t        |      i d}t        j                  | |dS )z
    Build an error response message.

    Args:
        request_id: Request identifier this responds to
        error: DirtyError instance or dict with error info

    Returns:
        dict: Error response message
    r   )
DirtyError)
error_typer   details)r]   r^   r	   )
errorsre   
isinstanceto_dictrP   r]   rE   strr   rL   )rW   r	   re   
error_dicts       r    make_error_responserm     sf     #%$]]_
	E4	 
 u+..5z

 ,, r=   c                 *    t         j                  | |dS )z
    Build a chunk message for streaming responses.

    Args:
        request_id: Request identifier this chunk belongs to
        data: Chunk data (must be JSON-serializable)

    Returns:
        dict: Chunk message
    )r]   r^   r!   )r   rM   )rW   r!   s     r    make_chunk_messagero   =  s     ,, r=   c                 (    t         j                  | dS )z
    Build an end-of-stream message.

    Args:
        request_id: Request identifier this ends

    Returns:
        dict: End message
    )r]   r^   )r   rN   )rW   s    r    make_end_messagerq   O  s     ** r=   )NN)rH   r/   r   r   rT   rh   r   r   rk   tuplerP   r`   rc   rm   ro   rq   rV   r=   r    <module>rs      s   
     &U Ut 59S C  -1=A2c d $C 4 <3  $  r=   