a
    còfœ  ã                   @   s:  d Z ddlZddlZddlZddlmZ zddlZdZW n e	yN   dZY n0 zddl
ZdZW n e	yv   dZY n0 ddlmZmZmZmZ ddlmZmZmZmZmZmZ dd	lmZmZ dd
lmZ ddlmZm Z m!Z!m"Z"m#Z# dZ$e %d¡Z&G dd„ de'ƒZ(G dd„ de'ƒZ)G dd„ de'ƒZ*G dd„ de'ƒZ+dS )z4Implementation of the X protocol for MySQL servers.
é    N)ÚBytesIOTFé   )ÚInterfaceErrorÚNotSupportedErrorÚOperationalErrorÚProgrammingError)Ú
ExprParserÚ
build_exprÚbuild_scalarÚbuild_bool_scalarÚbuild_int_scalarÚbuild_unsigned_int_scalar)Úencode_to_bytesÚget_item_or_attr)ÚColumn)ÚCRUD_PREPARE_MAPPINGÚSERVER_MESSAGESÚPROTOBUF_REPEATED_TYPESÚMessageÚmysqlxpb_enumiè  Zmysqlxc                   @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	Ú
CompressorzËImplements compression/decompression using `zstd_stream`, `lz4_message`
    and `deflate_stream` algorithms.

    Args:
        algorithm (str): Compression algorithm.

    .. versionadded:: 8.0.21

    c                 C   sR   || _ |dkr$t ¡ | _t ¡ | _n*|dkrBt ¡ | _t ¡ | _nd | _d | _d S )NÚzstd_streamZdeflate_stream)	Ú
_algorithmÚzstdZZstdCompressorÚ_compressobjZZstdDecompressorÚ_decompressobjÚzlibÚcompressobjÚdecompressobj©ÚselfÚ	algorithm© r"   úL/home/httpd/docs/test/DocsMgr/lib/python3.9/site-packages/mysqlx/protocol.pyÚ__init__J   s    

zCompressor.__init__c                 C   s’   | j dkr| j |¡S | j dkrptj ¡ �2}| ¡ }|| |¡7 }|| ¡ 7 }W d  ƒ n1 sb0    Y  |S | j |¡}|| j tj	¡7 }|S )z´Compresses data and returns it.

        Args:
            data (str, bytes or buffer object): Data to be compressed.

        Returns:
            bytes: Compressed data.
        r   Úlz4_messageN)
r   r   ÚcompressÚlz4ÚframeZLZ4FrameCompressorÚbeginÚflushr   ÚZ_SYNC_FLUSH)r    ÚdataZ
compressorÚ
compressedr"   r"   r#   r&   V   s    	

*zCompressor.compressc                 C   sz   | j dkr| j |¡S | j dkrXtj ¡ �}| |¡}W d  ƒ n1 sJ0    Y  |S | j |¡}|| j tj¡7 }|S )zÙDecompresses a frame of data and returns it as a string of bytes.

        Args:
            data (str, bytes or buffer object): Data to be compressed.

        Returns:
            bytes: Decompresssed data.
        r   r%   N)	r   r   Ú
decompressr'   r(   ZLZ4FrameDecompressorr*   r   r+   )r    r,   ZdecompressorÚdecompressedr"   r"   r#   r.   m   s    	

(zCompressor.decompressN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r$   r&   r.   r"   r"   r"   r#   r   @   s   	r   c                   @   s8   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
d„ ZdS )ÚMessageReaderz™Implements a Message Reader.

    Args:
        socket_stream (mysqlx.connection.SocketStream): `SocketStream` object.

    .. versionadded:: 8.0.21
    c                 C   s   || _ d | _d | _g | _d S ©N)Ú_streamÚ_compressorÚ_msgÚ
_msg_queue©r    Zsocket_streamr"   r"   r#   r$   ‹   s    zMessageReader.__init__c                 C   s  | j r| j  d¡S t d| j d¡¡\}}|dkr:tdƒ‚| j |d ¡}|tvr`td 	|¡ƒ‚|dkrx|d	krx|  
¡ S t ||¡}|d
k�r|d }t| j |d ¡ƒ}d}||k rüt d| d¡¡\}}	| |d ¡}
| j  t |	|
¡¡ ||d 7 }q®| j �r| j  d¡S dS |S )a¦  Reads X Protocol messages from the stream and returns a
        :class:`mysqlx.protobuf.Message` object.

        Raises:
            :class:`mysqlx.ProgrammingError`: If e connected server does not
                                              have the MySQL X protocol plugin
                                              enabled.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        r   ú<LBé   é
   z[The connected server does not have the MySQL X protocol plugin enabled or protocol mismatchr   zUnknown message type: {}é   ó    é   Úuncompressed_sizeÚpayloadé   N)r9   ÚpopÚstructÚunpackr6   Úreadr   r   Ú
ValueErrorÚformatÚ_read_messager   Zfrom_server_messager   r7   r.   Úappend)r    Ú
frame_sizeZ
frame_typeZframe_payloadZ	frame_msgrA   ÚstreamZbytes_processedZpayload_sizeÚmsg_typerB   r"   r"   r#   rJ   ‘   s2    
ÿ
ÿzMessageReader._read_messagec                 C   s"   | j dur| j }d| _ |S |  ¡ S )zgRead message.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        N)r8   rJ   ©r    Úmsgr"   r"   r#   Úread_messageÀ   s
    
zMessageReader.read_messagec                 C   s   | j durtdƒ‚|| _ dS )zÇPush message.

        Args:
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.

        Raises:
            :class:`mysqlx.OperationalError`: If message push slot is full.
        NzMessage push slot is full)r8   r   rO   r"   r"   r#   Úpush_messageÌ   s    	
zMessageReader.push_messagec                 C   s   |rt |ƒnd| _dS )zÏCreates a :class:`mysqlx.protocol.Compressor` object based on the
        compression algorithm.

        Args:
            algorithm (str): Compression algorithm.

        .. versionadded:: 8.0.21

        N©r   r7   r   r"   r"   r#   Úset_compressionÙ   s    
zMessageReader.set_compressionN)	r0   r1   r2   r3   r$   rJ   rQ   rR   rT   r"   r"   r"   r#   r4   ƒ   s   /r4   c                   @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	ÚMessageWriterzšImplements a Message Writer.

    Args:
        socket_stream (mysqlx.connection.SocketStream): `SocketStream` object.

    .. versionadded:: 8.0.21

    c                 C   s   || _ d | _d S r5   )r6   r7   r:   r"   r"   r#   r$   ï   s    zMessageWriter.__init__c                 C   s  |  |¡}| jrÔ|tkrÔt| ¡ ƒ}t d|d |¡}| j d ||g¡¡}t	dƒ}||d< |d |d< t	dƒ}||d< d t| 
¡ ƒd	d
… t| 
¡ ƒg¡}	tdƒ}
t dt|	ƒd |
¡}| j d ||	g¡¡ n4t| ¡ ƒ}t d|d |¡}| j d ||g¡¡ d	S )z™Write message.

        Args:
            msg_type (int): The message type.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.
        r;   r   r?   zMysqlx.Connection.CompressionZclient_messagesr<   rA   rB   Néþÿÿÿz&Mysqlx.ClientMessages.Type.COMPRESSION)Z	byte_sizer7   Ú_COMPRESSION_THRESHOLDr   Zserialize_to_stringrE   Úpackr&   Újoinr   Zserialize_partial_to_stringr   Úlenr6   Úsendall)r    rN   rP   Zmsg_sizeZmsg_strÚheaderr-   Zmsg_first_fieldsZmsg_payloadÚoutputZmsg_comp_idr"   r"   r#   Úwrite_messageó   s*    

þÿzMessageWriter.write_messagec                 C   s   |rt |ƒnd| _dS )z¬Creates a :class:`mysqlx.protocol.Compressor` object based on the
        compression algorithm.

        Args:
            algorithm (str): Compression algorithm.
        NrS   r   r"   r"   r#   rT     s    zMessageWriter.set_compressionN)r0   r1   r2   r3   r$   r^   rT   r"   r"   r"   r#   rU   æ   s   "rU   c                   @   s  e Zd ZdZdd„ Zedd„ ƒZdd„ Zdd	„ ZdDdd„Z	dd„ Z
dd„ Zdd„ Zdd„ Zdd„ ZdEd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d*d+„ Zd,d-„ Zd.d/„ ZdFd0d1„Zd2d3„ Zd4d5„ Zd6d7„ Zd8d9„ Zd:d;„ Z d<d=„ Z!d>d?„ Z"d@dA„ Z#dGdBdC„Z$dS )HÚProtocolzàImplements the MySQL X Protocol.

    Args:
        read (mysqlx.protocol.MessageReader): A Message Reader object.
        writer (mysqlx.protocol.MessageWriter): A Message Writer object.

    .. versionchanged:: 8.0.21
    c                 C   s   || _ || _d | _g | _d S r5   )Ú_readerÚ_writerÚ_compression_algorithmÚ	_warnings)r    ÚreaderÚwriterr"   r"   r#   r$   (  s    zProtocol.__init__c                 C   s   | j S )z'str: The compresion algorithm.
        )rb   )r    r"   r"   r#   Úcompression_algorithm.  s    zProtocol.compression_algorithmc                 C   sX   |j r| ¡ |d< |jr*|d  | ¡ ¡ |jrB|d  | ¡ ¡ |jrT| ¡ |d< dS )z­Apply filter.

        Args:
            msg (mysqlx.protobuf.Message): The MySQL X Protobuf Message.
            stmt (Statement): A `Statement` based type object.
        ÚcriteriaÚorderÚgroupingZgrouping_criteriaN)	Z	has_whereZget_where_exprZhas_sortÚextendZget_sort_exprZhas_group_byZget_groupingZ
has_havingZ
get_having)r    rP   Ústmtr"   r"   r#   Ú_apply_filter4  s    zProtocol._apply_filterc                 C   sö  t |tƒr2td|d�}tdd|d�}tdd|d�S t |tƒrNtddt|ƒd�S t |tƒr„|d	k rrtddt|ƒd�S tddt|ƒd�S t |tƒrÖt	|ƒd
krÖ|\}}td||  
|¡d�}td| ¡ gd�}tdd
|d�S t |tƒ�st |ttfƒ�r’t |d	 tƒ�r’g }|D ]h}	g }
|	 ¡ D ],\}}td||  
|¡d�}|
 | ¡ ¡ �qtd|
d�}tdd
|d�}| | ¡ ¡ �q
tdƒ}||d< tdd|d�S t |tƒ�ròg }
|D ],\}}td||  
|¡d�}|
 | ¡ ¡ �q¦td|
d�}tdd
|d�}|S dS )z Create any.

        Args:
            arg (object): Arbitrary object.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        zMysqlx.Datatypes.Scalar.String©ÚvaluezMysqlx.Datatypes.Scalaré   )ÚtypeZv_stringúMysqlx.Datatypes.Anyr   )rp   Úscalarr   é   ú#Mysqlx.Datatypes.Object.ObjectField©Úkeyrn   úMysqlx.Datatypes.Object©Zfld©rp   ÚobjzMysqlx.Datatypes.Arrayrn   é   )rp   ÚarrayN)Ú
isinstanceÚstrr   Úboolr   Úintr   r   ÚtuplerZ   Ú_create_anyÚget_messageÚdictÚlistÚitemsrK   )r    Úargrn   rr   Zarg_keyÚ	arg_valueÚobj_fldrz   Zarray_valuesr†   Úobj_fldsrv   Úmsg_objÚmsg_anyrP   r"   r"   r#   r‚   D  sj    	

ÿ
ÿÿ
ÿÿÿ
ÿ
ÿzProtocol._create_anyTc           
         sž   ‡‡fdd„‰ |  ¡ }| ¡ }|du r8‡ fdd„|D ƒS t|ƒ}|dg }|t|ƒkr^tdƒ‚| ¡ D ]2\}}||vr„td |¡ƒ‚|| }	ˆ |ƒ||	< qf|S )a›  Returns the binding any/scalar.

        Args:
            stmt (Statement): A `Statement` based type object.
            is_scalar (bool): `True` to return scalar values.

        Raises:
            :class:`mysqlx.ProgrammingError`: If unable to find placeholder for
                                              parameter.

        Returns:
            list: A list of ``Any`` or ``Scalar`` objects.
        c                    s   ˆ rt | ƒ ¡ S ˆ | ¡ ¡ S r5   )r
   rƒ   r‚   rm   )Ú	is_scalarr    r"   r#   Ú<lambda>Œ  s    ÿz,Protocol._get_binding_args.<locals>.<lambda>Nc                    s   g | ]}ˆ |ƒ‘qS r"   r"   ©Ú.0rn   )Úbuild_valuer"   r#   Ú
<listcomp>“  r?   z.Protocol._get_binding_args.<locals>.<listcomp>z;The number of bind parameters and placeholders do not matchz-Unable to find placeholder for parameter: {0})Úget_bindingsZget_binding_maprZ   r   r†   rI   )
r    rk   r�   ZbindingsZbinding_mapÚcountÚargsÚnamern   Úposr"   )r‘   r�   r    r#   Ú_get_binding_args~  s"    
ÿzProtocol._get_binding_argsc                 C   s(  |d dkrRt  d|d ¡}| j |j¡ t d|j|j¡ | |j	|j|j¡ nÒ|d dkrpt  d|d ¡ n´|d dk�r$t  d	|d ¡}|d
 t
dƒkr¸| dd„ |d D ƒ¡ nlt|d ttƒƒrÖ|d d n|d }|d
 t
dƒk�r| t|dƒ¡ n"|d
 t
dƒk�r$| t|dƒ¡ dS )z¨Process frame.

        Args:
            msg (mysqlx.protobuf.Message): A MySQL X Protobuf Message.
            result (Result): A `Result` based type object.
        rp   r   zMysqlx.Notice.WarningrB   z:Protocol.process_frame Received Warning Notice code %s: %srs   z$Mysqlx.Notice.SessionVariableChangedr{   z!Mysqlx.Notice.SessionStateChangedÚparamzBMysqlx.Notice.SessionStateChanged.Parameter.GENERATED_DOCUMENT_IDSc                 S   s    g | ]}t t |d ƒdƒ ¡ ‘qS )Zv_octetsrn   )r   Údecoder�   r"   r"   r#   r’   º  s   þ
ÿz+Protocol._process_frame.<locals>.<listcomp>rn   r   z9Mysqlx.Notice.SessionStateChanged.Parameter.ROWS_AFFECTEDZv_unsigned_intz?Mysqlx.Notice.SessionStateChanged.Parameter.GENERATED_INSERT_IDN)r   Zfrom_messagerc   rK   rP   Ú_LOGGERÚwarningÚcodeZappend_warningÚlevelr   Zset_generated_idsr}   r�   r   Zset_rows_affectedr   Zset_generated_insert_id)r    rP   ÚresultZwarn_msgZsess_state_msgZsess_state_valuer"   r"   r#   Ú_process_frame¢  sV    ÿÿÿÿÿþÿÿÿýÿÿÿÿzProtocol._process_framec              
   C   sú   z| j  ¡ }W nD tyR } z,t| ¡ ƒ}|r>td ||¡ƒ‚W Y d}~n
d}~0 0 |jdkrrt|d |d ƒ‚q |jdkr z|  ||¡ W qô   Y q Y qô0 q |jdkr®dS |jdkrÄ| 	d	¡ q |jd
krÚ| 
d	¡ q |jdkrö| d	¡ qöq qöq |S )z`Read message.

        Args:
            result (Result): A `Result` based type object.
        z{} reason: {}NúMysqlx.ErrorrP   r�   úMysqlx.Notice.FramezMysqlx.Sql.StmtExecuteOkzMysqlx.Resultset.FetchDoneTz(Mysqlx.Resultset.FetchDoneMoreResultsetsúMysqlx.Resultset.Row)r`   rQ   ÚRuntimeErrorÚreprZget_warningsrI   rp   r   r    Z
set_closedZset_has_more_resultsZset_has_data)r    rŸ   rP   ÚerrÚwarningsr"   r"   r#   rJ   Í  s4    
ÿ






zProtocol._read_messagec                 C   s"   || _ | j |¡ | j |¡ dS )zðSets the compression algorithm to be used by the compression
        object, for uplink and downlink.

        Args:
            algorithm (str): Algorithm to be used in compression/decompression.

        .. versionadded:: 8.0.21

        N)rb   r`   rT   ra   r   r"   r"   r#   rT   ï  s    
zProtocol.set_compressionc                 C   sZ   t dƒ}| j tdƒ|¡ | j ¡ }|jdkr:| j ¡ }q$|jdkrVt|d |d ƒ‚|S )zkGet capabilities.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        z!Mysqlx.Connection.CapabilitiesGetz/Mysqlx.ClientMessages.Type.CON_CAPABILITIES_GETr¢   r¡   rP   r�   )r   ra   r^   r   r`   rQ   rp   r   rO   r"   r"   r#   Úget_capabilitesý  s    þ


zProtocol.get_capabilitesc              
   K   s$  |sdS t dƒ}| ¡ D ]¤\}}t dƒ}||d< t|tƒrš|}g }|D ]*}t d||  || ¡d�}	| |	 ¡ ¡ qFt d|d�}
t d	d
|
d�}| ¡ |d< n|  |¡|d< |d  | ¡ g¡ qt dƒ}||d< | j 	t
dƒ|¡ z
|  ¡ W S  t�y } z|jdk�r
‚ W Y d}~n
d}~0 0 dS )z­Set capabilities.

        Args:
            **kwargs: Arbitrary keyword arguments.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        NzMysqlx.Connection.CapabilitieszMysqlx.Connection.Capabilityr–   rt   ru   rw   rx   rq   rs   ry   rn   Úcapabilitiesz!Mysqlx.Connection.CapabilitiesSetz/Mysqlx.ClientMessages.Type.CON_CAPABILITIES_SETiŠ  )r   r†   r}   r„   r‚   rK   rƒ   rj   ra   r^   r   Úread_okr   Úerrno)r    Úkwargsr©   rv   rn   Z
capabilityr†   rŠ   Úitemr‰   r‹   rŒ   rP   r¦   r"   r"   r#   Úset_capabilities  s@    	
þþ
zProtocol.set_capabilitiesNc                 C   sF   t dƒ}||d< |dur ||d< |dur0||d< | j tdƒ|¡ dS )zÖSend authenticate start.

        Args:
            method (str): Message method.
            auth_data (Optional[str]): Authentication data.
            initial_response (Optional[str]): Initial response.
        z Mysqlx.Session.AuthenticateStartZ	mech_nameNÚ	auth_dataÚinitial_responsez2Mysqlx.ClientMessages.Type.SESS_AUTHENTICATE_START©r   ra   r^   r   )r    Úmethodr¯   r°   rP   r"   r"   r#   Úsend_auth_start=  s    ÿÿzProtocol.send_auth_startc                 C   s:   | j  ¡ }|jdkr | j  ¡ }q
|jdkr2tdƒ‚|d S )züRead authenticate continue.

        Raises:
            :class:`InterfaceError`: If the message type is not
                                     `Mysqlx.Session.AuthenticateContinue`

        Returns:
            str: The authentication data.
        r¢   ú#Mysqlx.Session.AuthenticateContinuez>Unexpected message encountered during authentication handshaker¯   )r`   rQ   rp   r   rO   r"   r"   r#   Úread_auth_continueN  s    



zProtocol.read_auth_continuec                 C   s"   t d|d�}| j tdƒ|¡ dS )zeSend authenticate continue.

        Args:
            auth_data (str): Authentication data.
        r´   )r¯   z5Mysqlx.ClientMessages.Type.SESS_AUTHENTICATE_CONTINUENr±   )r    r¯   rP   r"   r"   r#   Úsend_auth_continue`  s    ÿÿÿzProtocol.send_auth_continuec                 C   s0   | j  ¡ }|jdkrq,|jdkr t|jƒ‚q dS )z~Read authenticate OK.

        Raises:
            :class:`mysqlx.InterfaceError`: If message type is `Mysqlx.Error`.
        zMysqlx.Session.AuthenticateOkr¡   N)r`   rQ   rp   r   rP   rO   r"   r"   r#   Úread_auth_okk  s
    


zProtocol.read_auth_okc                 C   s@  |j rÂ|jdkrÂ|jdkr*|  |¡\}}nB|jdkrD|  |¡\}}n(|jdkr^|  |¡\}}ntd |¡ƒ‚t| ¡ ƒ}t	dƒ}t
dƒ}t
d||d	�|d
< |jdkrºt
d||d d	�|d< ||d< t| \}}	t
dƒ}
t	|ƒ|
d< ||
|	< t
dƒ}|j|d< |
|d< | j t	dƒ|¡ z|  ¡  W n t�y:   t‚Y n0 dS )a¦  
        Send prepare statement.

        Args:
            msg_type (str): Message ID string.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.
            stmt (Statement): A `Statement` based type object.

        Raises:
            :class:`mysqlx.NotSupportedError`: If prepared statements are not
                                               supported.

        .. versionadded:: 8.0.16
        úMysqlx.Crud.InsertúMysqlx.Crud.FindúMysqlx.Crud.UpdateúMysqlx.Crud.DeletezInvalid message type: {}z!Mysqlx.Expr.Expr.Type.PLACEHOLDERzMysqlx.Crud.LimitExprzMysqlx.Expr.Expr)rp   ÚpositionÚ	row_countr   ÚoffsetZ
limit_exprú#Mysqlx.Prepare.Prepare.OneOfMessagerp   zMysqlx.Prepare.PrepareÚstmt_idrk   z*Mysqlx.ClientMessages.Type.PREPARE_PREPAREN)Ú	has_limitrp   Ú
build_findÚbuild_updateÚbuild_deleterH   rI   rZ   r“   r   r   r   rÀ   ra   r^   rª   r   r   )r    rN   rP   rk   Ú_r¼   ÚplaceholderZmsg_limit_exprÚ
oneof_typeÚoneof_opÚ	msg_oneofZmsg_preparer"   r"   r#   Úsend_prepare_preparex  sH    


þ

þ

þzProtocol.send_prepare_preparec           	      C   s¤   t | \}}tdƒ}t|ƒ|d< |||< tdƒ}|j|d< | j|dd�}|rZ|d  |¡ |jrŽ|d  |  | ¡ ¡ 	¡ |  | 
¡ ¡ 	¡ g¡ | j tdƒ|¡ d	S )
a  
        Send execute statement.

        Args:
            msg_type (str): Message ID string.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.
            stmt (Statement): A `Statement` based type object.

        .. versionadded:: 8.0.16
        r¿   rp   zMysqlx.Prepare.ExecuterÀ   F©r�   r•   z*Mysqlx.ClientMessages.Type.PREPARE_EXECUTEN)r   r   r   rÀ   r˜   rj   rÁ   r‚   Úget_limit_row_countrƒ   Úget_limit_offsetra   r^   )	r    rN   rP   rk   rÇ   rÈ   rÉ   Zmsg_executer•   r"   r"   r#   Úsend_prepare_execute¯  s$    
þþzProtocol.send_prepare_executec                 C   s.   t dƒ}||d< | j tdƒ|¡ |  ¡  dS )zŽ
        Send prepare deallocate statement.

        Args:
            stmt_id (int): Statement ID.

        .. versionadded:: 8.0.16
        zMysqlx.Prepare.DeallocaterÀ   z-Mysqlx.ClientMessages.Type.PREPARE_DEALLOCATEN)r   ra   r^   r   rª   )r    rÀ   Zmsg_deallocr"   r"   r#   Úsend_prepare_deallocateÏ  s    	þz Protocol.send_prepare_deallocatec                 C   sx   |j r8tdƒ}| ¡ |d< |jdkr0| ¡ |d< ||d< |dkrDdnd}| j||d	�}|rh|d
  |¡ |  ||¡ dS )a)  
        Send a message without prepared statements support.

        Args:
            msg_type (str): Message ID string.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.
            stmt (Statement): A `Statement` based type object.

        .. versionadded:: 8.0.16
        zMysqlx.Crud.Limitr½   r¹   r¾   Úlimitú+Mysqlx.ClientMessages.Type.SQL_STMT_EXECUTEFTrË   r•   N)rÁ   r   rÌ   rp   rÍ   r˜   rj   Úsend_msg)r    rN   rP   rk   Z	msg_limitr�   r•   r"   r"   r#   Úsend_msg_without_psß  s    
ÿþzProtocol.send_msg_without_psc                 C   s   | j  t|ƒ|¡ dS )zÆ
        Send a message.

        Args:
            msg_type (str): Message ID string.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.

        .. versionadded:: 8.0.16
        N)ra   r^   r   )r    rN   rP   r"   r"   r#   rÒ   ø  s    
zProtocol.send_msgc                 C   sœ   t | ¡ rdndƒ}td|jj|jjd�}td||d�}|jrJ| ¡ |d< |  ||¡ | 	¡ rlt dƒ|d	< n| 
¡ r€t d
ƒ|d	< |jdkr”|j|d< d|fS )a‹  Build find/read message.

        Args:
            stmt (Statement): A :class:`mysqlx.ReadStatement` or
                              :class:`mysqlx.FindStatement` object.

        Returns:
            (tuple): Tuple containing:

                * `str`: Message ID string.
                * :class:`mysqlx.protobuf.Message`: MySQL X Protobuf Message.

        .. versionadded:: 8.0.16
        úMysqlx.Crud.DataModel.DOCUMENTúMysqlx.Crud.DataModel.TABLEúMysqlx.Crud.Collection©r–   Úschemar¹   ©Ú
data_modelÚ
collectionÚ
projectionz'Mysqlx.Crud.Find.RowLock.EXCLUSIVE_LOCKZlockingz$Mysqlx.Crud.Find.RowLock.SHARED_LOCKr   Zlocking_optionsz$Mysqlx.ClientMessages.Type.CRUD_FIND)r   Úis_doc_basedr   Útargetr–   rØ   Zhas_projectionZget_projection_exprrl   Zis_lock_exclusiveZis_lock_sharedZlock_contention©r    rk   rÚ   rÛ   rP   r"   r"   r#   rÂ     s0    ÿþþÿÿÿ

zProtocol.build_findc                 C   sª   t | ¡ rdndƒ}td|jj|jjd�}td||d�}|  ||¡ | ¡  ¡ D ]P\}}tdƒ}|j	|d< |j
|d	< |jd
urŒt|jƒ|d< |d  | ¡ g¡ qPd|fS )aŒ  Build update message.

        Args:
            stmt (Statement): A :class:`mysqlx.ModifyStatement` or
                              :class:`mysqlx.UpdateStatement` object.

        Returns:
            (tuple): Tuple containing:

                * `str`: Message ID string.
                * :class:`mysqlx.protobuf.Message`: MySQL X Protobuf Message.

        .. versionadded:: 8.0.16
        rÔ   rÕ   rÖ   r×   rº   rÙ   zMysqlx.Crud.UpdateOperationÚ	operationÚsourceNrn   z&Mysqlx.ClientMessages.Type.CRUD_UPDATE)r   rÝ   r   rÞ   r–   rØ   rl   Zget_update_opsr†   Zupdate_typerá   rn   r	   rj   rƒ   )r    rk   rÚ   rÛ   rP   rÅ   Z	update_oprà   r"   r"   r#   rÃ   +  s*    ÿþþÿ


zProtocol.build_updatec                 C   sL   t | ¡ rdndƒ}td|jj|jjd�}td||d�}|  ||¡ d|fS )aŒ  Build delete message.

        Args:
            stmt (Statement): A :class:`mysqlx.DeleteStatement` or
                              :class:`mysqlx.RemoveStatement` object.

        Returns:
            (tuple): Tuple containing:

                * `str`: Message ID string.
                * :class:`mysqlx.protobuf.Message`: MySQL X Protobuf Message.

        .. versionadded:: 8.0.16
        rÔ   rÕ   rÖ   r×   r»   rÙ   z&Mysqlx.ClientMessages.Type.CRUD_DELETE)r   rÝ   r   rÞ   r–   rØ   rl   rß   r"   r"   r#   rÄ   M  s    ÿþ
ÿÿzProtocol.build_deletec                 C   s|   t d||dd�}|rtg }| ¡ D ]*\}}t d||  |¡d�}| | ¡ ¡ q t d|d�}	t dd	|	d
�}
|
 ¡ g|d< d|fS )aª  Build execute statement.

        Args:
            namespace (str): The namespace.
            stmt (Statement): A `Statement` based type object.
            fields (Optional[dict]): The message fields.

        Returns:
            (tuple): Tuple containing:

                * `str`: Message ID string.
                * :class:`mysqlx.protobuf.Message`: MySQL X Protobuf Message.

        .. versionadded:: 8.0.16
        zMysqlx.Sql.StmtExecuteF)Ú	namespacerk   Zcompact_metadatart   ru   rw   rx   rq   rs   ry   r•   rÑ   )r   r†   r‚   rK   rƒ   )r    râ   rk   ÚfieldsrP   rŠ   rv   rn   r‰   r‹   rŒ   r"   r"   r#   Úbuild_execute_statementf  s    ÿ
ÿz Protocol.build_execute_statementc           
      C   s  t | ¡ rdndƒ}td|jj|jjd�}td||d�}t|dƒrv|jD ],}t|| ¡  ƒ 	¡ }|d  
| ¡ g¡ qH| ¡ D ]f}td	ƒ}t|tƒr¸|D ]}	|d
  
t|	ƒ ¡ g¡ q˜n|d
  
t|ƒ ¡ g¡ |d  
| ¡ g¡ q~t|dƒrü| ¡ |d< d|fS )a‹  Build insert statement.

        Args:
            stmt (Statement): A :class:`mysqlx.AddStatement` or
                              :class:`mysqlx.InsertStatement` object.

        Returns:
            (tuple): Tuple containing:

                * `str`: Message ID string.
                * :class:`mysqlx.protobuf.Message`: MySQL X Protobuf Message.

        .. versionadded:: 8.0.16
        rÔ   rÕ   rÖ   r×   r¸   rÙ   Ú_fieldsrÜ   zMysqlx.Crud.Insert.TypedRowÚfieldÚrowÚ	is_upsertZupsertz&Mysqlx.ClientMessages.Type.CRUD_INSERT)r   rÝ   r   rÞ   r–   rØ   Úhasattrrå   r   Zparse_table_insert_fieldrj   rƒ   Z
get_valuesr}   r…   r	   rè   )
r    rk   rÚ   rÛ   rP   ræ   Úexprrn   rç   Úvalr"   r"   r#   Úbuild_insert„  s4    ÿþþÿ



zProtocol.build_insertc                 C   s   |   |¡}|durtdƒ‚dS )z¼Close the result.

        Args:
            result (Result): A `Result` based type object.

        Raises:
            :class:`mysqlx.OperationalError`: If message read is None.
        NzExpected to close the result)rJ   r   ©r    rŸ   rP   r"   r"   r#   Úclose_result°  s    	
zProtocol.close_resultc                 C   s4   |   |¡}|du rdS |jdkr$|S | j |¡ dS )z\Read row.

        Args:
            result (Result): A `Result` based type object.
        Nr£   )rJ   rp   r`   rR   rí   r"   r"   r#   Úread_row½  s    

zProtocol.read_rowc                 C   s²   g }|   |¡}|du rq®|jdkr0| j |¡ q®|jdkrBtdƒ‚t|d |d |d |d |d	 |d
 |d | dd¡| dd¡| dd¡| dd¡| d¡ƒ}| |¡ q|S )z¿Returns column metadata.

        Args:
            result (Result): A `Result` based type object.

        Raises:
            :class:`mysqlx.InterfaceError`: If unexpected message.
        Nr£   zMysqlx.Resultset.ColumnMetaDatazUnexpected msg typerp   ÚcatalogrØ   ÚtableZoriginal_tabler–   Úoriginal_nameÚlengthé   Z	collationr   Zfractional_digitsÚflagsé   Úcontent_type)rJ   rp   r`   rR   r   r   ÚgetrK   )r    rŸ   ÚcolumnsrP   Úcolr"   r"   r#   Úget_column_metadataË  s(    	






ùzProtocol.get_column_metadatac                 C   sD   | j  ¡ }|jdkr.td |d ¡|d d�‚|jdkr@tdƒ‚dS )	zeRead OK.

        Raises:
            :class:`mysqlx.InterfaceError`: If unexpected message.
        r¡   zMysqlx.Error: {}rP   r�   )r«   z	Mysqlx.OkzUnexpected message encounteredN)r`   rQ   rp   r   rI   rO   r"   r"   r#   rª   é  s    

ÿ
zProtocol.read_okc                 C   s   t dƒ}| j tdƒ|¡ dS )zSend connection close.zMysqlx.Connection.Closez$Mysqlx.ClientMessages.Type.CON_CLOSENr±   rO   r"   r"   r#   Úsend_connection_closeö  s    ÿÿzProtocol.send_connection_closec                 C   s   t dƒ}| j tdƒ|¡ dS )zSend close.zMysqlx.Session.Closez%Mysqlx.ClientMessages.Type.SESS_CLOSENr±   rO   r"   r"   r#   Ú
send_closeü  s    ÿÿzProtocol.send_closec                 C   sL   t dƒ}tdƒ}||d< d|d< tdƒ}| ¡ g|d< | j t dƒ|¡ d	S )
zSend expectation.z3Mysqlx.Expect.Open.Condition.Key.EXPECT_FIELD_EXISTzMysqlx.Expect.Open.ConditionZcondition_keyz6.1Zcondition_valuezMysqlx.Expect.OpenZcondz&Mysqlx.ClientMessages.Type.EXPECT_OPENN)r   r   rƒ   ra   r^   )r    Zcond_keyZmsg_ocZmsg_eor"   r"   r#   Úsend_expect_open  s    ÿÿÿzProtocol.send_expect_openc                 C   sr   t dƒ}|du r@z|  ¡  |  ¡  d}W n ty>   d}Y n0 |rLd|d< | j tdƒ|¡ |  ¡  |rndS dS )z¨Send reset session message.

        Returns:
            boolean: ``True`` if the server will keep the session open,
                     otherwise ``False``.
        zMysqlx.Session.ResetNTFÚ	keep_openz%Mysqlx.ClientMessages.Type.SESS_RESET)r   rþ   rª   r   ra   r^   r   )r    rÿ   rP   r"   r"   r#   Ú
send_reset  s&    
ÿÿzProtocol.send_reset)T)NN)N)N)%r0   r1   r2   r3   r$   Úpropertyrf   rl   r‚   r˜   r    rJ   rT   r¨   r®   r³   rµ   r¶   r·   rÊ   rÎ   rÏ   rÓ   rÒ   rÂ   rÃ   rÄ   rä   rì   rî   rï   rû   rª   rü   rý   rþ   r   r"   r"   r"   r#   r_     sD   
:
$+"-
7 '"
,r_   ),r3   ÚloggingrE   r   Úior   Z	lz4.framer'   ZHAVE_LZ4ÚImportErrorZ	zstandardr   Z	HAVE_ZSTDÚerrorsr   r   r   r   rê   r   r	   r
   r   r   r   Úhelpersr   r   rŸ   r   Úprotobufr   r   r   r   r   rW   Ú	getLoggerr›   Úobjectr   r4   rU   r_   r"   r"   r"   r#   Ú<module>   s2   

 
Cc9