3
Šãde§  ã               @   s¨  d Z ddlZddlZddlmZ ddlmZmZmZm	Z	m
Z
mZ yddlZdZW n ek
rh   dZY nX yddlZdZW n ek
r’   dZY nX 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dl)m*Z* ddl+m,Z,m-Z-m.Z.m/Z/m0Z0m1Z1m2Z2m3Z3m4Z4m5Z5 ddl6m7Z7m8Z8m9Z9m:Z:m;Z;m<Z<m=Z=m>Z> dZ?G dd„ dƒZ@G dd„ dƒZAG dd„ dƒZBG dd„ dƒZCdS )z3Implementation of the X protocol for MySQL servers.é    N)ÚBytesIO)ÚAnyÚDictÚListÚOptionalÚTupleÚUnionTFé   )ÚInterfaceErrorÚNotSupportedErrorÚOperationalErrorÚProgrammingError)Ú
ExprParserÚbuild_bool_scalarÚ
build_exprÚbuild_int_scalarÚbuild_scalarÚbuild_unsigned_int_scalar)Úencode_to_bytesÚget_item_or_attr)Úlogger)ÚCRUD_PREPARE_MAPPINGÚPROTOBUF_REPEATED_TYPESÚSERVER_MESSAGESÚMessageÚmysqlxpb_enum)ÚColumn)
ÚAddStatementÚDeleteStatementÚFilterableStatementÚFindStatementÚInsertStatementÚModifyStatementÚReadStatementÚRemoveStatementÚSqlStatementÚUpdateStatement)Ú
ColumnTypeÚMessageTypeÚProtobufMessageCextTypeÚProtobufMessageTypeÚResultBaseTypeÚ
SocketTypeÚStatementTypeÚ
StrOrBytesiè  c               @   s@   e Zd ZdZeddœdd„Zeedœdd„Zeedœ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

    N)Ú	algorithmÚreturnc             C   sP   || _ d | _d | _|dkr0tjƒ | _tjƒ | _n|dkrLtjƒ | _tjƒ | _d S )NÚzstd_streamZdeflate_stream)	Ú
_algorithmÚ_compressobjÚ_decompressobjÚzstdZZstdCompressorZZstdDecompressorÚzlibÚcompressobjÚdecompressobj)Úselfr0   © r;   úF/var/www/agendate/envp3/lib/python3.6/site-packages/mysqlx/protocol.pyÚ__init__p   s    

zCompressor.__init__)Údatar1   c          
   C   s~   | j dkr| jj|ƒS | j dkr\tjjƒ �(}|jƒ }||j|ƒ7 }||jƒ 7 }W dQ R X |S | jj|ƒ}|| 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.
        r2   Úlz4_messageN)
r3   r4   ÚcompressÚlz4ÚframeZLZ4FrameCompressorÚbeginÚflushr7   ÚZ_SYNC_FLUSH)r:   r>   Z
compressorÚ
compressedr;   r;   r<   r@   |   s    	

zCompressor.compressc          
   C   sf   | j dkr| jj|ƒS | j dkrDtjjƒ �}|j|ƒ}W dQ R X |S | jj|ƒ}|| 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.
        r2   r?   N)	r3   r5   Ú
decompressrA   rB   ZLZ4FrameDecompressorrD   r7   rE   )r:   r>   ZdecompressorÚdecompressedr;   r;   r<   rG   “   s    	

zCompressor.decompress)
Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ústrr=   r.   Úbytesr@   rG   r;   r;   r;   r<   r/   e   s   	r/   c               @   s\   e Zd ZdZeddœdd„Zedœdd„Zedœd	d
„Zeddœdd„Z	e
ddœdd„ZdS )ÚMessageReaderz™Implements a Message Reader.

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

    .. versionadded:: 8.0.21
    N)Úsocket_streamr1   c             C   s   || _ d | _d | _g | _d S )N)Ú_streamÚ_compressorÚ_msgÚ
_msg_queue)r:   rP   r;   r;   r<   r=   ²   s    zMessageReader.__init__)r1   c             C   s  | j r| j jdƒS tjd| jjdƒƒ\}}|dkr:tdƒ‚| jj|d ƒ}|tkr`td|› �ƒ‚|dkrx|d	krx| j	ƒ S t
j||ƒ}|d
k�r|d }t| jj|d ƒƒ}d}xP||k rþtjd|jdƒƒ\}}	|j|d ƒ}
| j jt
j|	|
ƒƒ ||d 7 }q°W | j �r| j 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   z<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)rT   ÚpopÚstructÚunpackrQ   Úreadr   r   Ú
ValueErrorÚ_read_messager   Zfrom_server_messager   rR   rG   Úappend)r:   Ú
frame_sizeZ
frame_typeZframe_payloadZ	frame_msgrZ   ÚstreamZbytes_processedZpayload_sizeÚmsg_typer[   r;   r;   r<   rb   ¸   s.    

zMessageReader._read_messagec             C   s"   | j dk	r| j }d| _ |S | jƒ S )zgRead message.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        N)rS   rb   )r:   Úmsgr;   r;   r<   Úread_messageç   s
    
zMessageReader.read_message)rg   r1   c             C   s   | j dk	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)rS   r   )r:   rg   r;   r;   r<   Úpush_messageó   s    	
zMessageReader.push_message)r0   r1   c             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/   rR   )r:   r0   r;   r;   r<   Úset_compression   s    
zMessageReader.set_compression)rI   rJ   rK   rL   r,   r=   r(   rb   rh   ri   rM   rj   r;   r;   r;   r<   rO   ©   s   /rO   c               @   sB   e Zd ZdZeddœdd„Zeeddœdd„Ze	dd	œd
d„Z
dS )ÚMessageWriterzšImplements a Message Writer.

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

    .. versionadded:: 8.0.21

    N)rP   r1   c             C   s   || _ d | _d S )N)rQ   rR   )r:   rP   r;   r;   r<   r=     s    zMessageWriter.__init__)rf   rg   r1   c             C   s  |j |ƒ}| jrÔ|tkrÔt|jƒ ƒ}tjd|d |ƒ}| jjdj||gƒƒ}t	dƒ}||d< |d |d< t	dƒ}||d< djt|j
ƒ ƒd	d… t|j
ƒ ƒgƒ}	tdƒ}
tjdt|	ƒd |
ƒ}| jjdj||	gƒƒ n4t|jƒ ƒ}tjd|d |ƒ}| jjdj||gƒƒ d	S )z™Write message.

        Args:
            msg_type (int): The message type.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.
        z<LBr	   rX   zMysqlx.Connection.CompressionZclient_messagesrU   rZ   r[   Né   z&Mysqlx.ClientMessages.Type.COMPRESSIONéþÿÿÿ)Z	byte_sizerR   Ú_COMPRESSION_THRESHOLDr   Zserialize_to_stringr^   Úpackr@   Újoinr   Zserialize_partial_to_stringr   ÚlenrQ   Úsendall)r:   rf   rg   Zmsg_sizeZmsg_strÚheaderrF   Zmsg_first_fieldsZmsg_payloadÚoutputZmsg_comp_idr;   r;   r<   Úwrite_message  s(    
zMessageWriter.write_message)r0   r1   c             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.
        N)r/   rR   )r:   r0   r;   r;   r<   rj   @  s    zMessageWriter.set_compression)rI   rJ   rK   rL   r,   r=   Úintr(   ru   rM   rj   r;   r;   r;   r<   rk     s   %rk   c            	   @   sæ  e Zd ZdZeeddœdd„Zeee	 dœdd„ƒZ
eeedd	œd
d„ƒZeee dœdd„ZdUeeef eeed eeeef  f dœdd„Zeeddœdd„Zeee dœdd„Ze	ddœdd„Zedœdd„Zeddœdd „ZdVe	ee	 ee	 dd!œd"d#„Zedœd$d%„Z e	dd&œd'd(„Z!ddœd)d*„Z"e	eee#e$e%e&e'e(f dd+œd,d-„Z)e	eedd+œd.d/„Z*e+dd0œd1d2„Z,e	eeeef dd+œd3d4„Z-e	edd5œd6d7„Z.ee#e&f e/e	ef d8œd9d:„Z0ee%e(f e/e	ef d8œd;d<„Z1ee$e'f e/e	ef d8œd=d>„Z2dWe	ee	e3f ee4e	ef  e/e	ef d?œd@dA„Z5eee6e7f e/e	ef d8œdBdC„ƒZ8eddœdDdE„Z9eee dœdFdG„Z:eee; dœdHdI„Z<ddœdJdK„Z=ddœdLdM„Z>ddœdNdO„Z?ddœdPdQ„Z@dXee edRœdSdT„ZAdS )YÚ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
    N)ÚreaderÚwriterr1   c             C   s   || _ || _d | _g | _d S )N)Ú_readerÚ_writerÚ_compression_algorithmÚ	_warnings)r:   rx   ry   r;   r;   r<   r=   T  s    zProtocol.__init__)r1   c             C   s   | j S )zstr: The compresion algorithm.)r|   )r:   r;   r;   r<   Úcompression_algorithmZ  s    zProtocol.compression_algorithm)rg   Ústmtr1   c             C   sX   |j r|jƒ | d< |jr*| d j|jƒ ƒ |jrB| d j|jƒ ƒ |jrT|jƒ | 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)rg   r   r;   r;   r<   Ú_apply_filter_  s    zProtocol._apply_filter)Úargr1   c             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|| j
|ƒd�}td|jƒ gd�}tdd
|d�S t |tƒsþt |ttfƒoút |d	 tƒ�r–g }xt|D ]l}	g }
x8|	jƒ D ],\}}td|| j
|ƒd�}|
j|jƒ ƒ �qW td|
d�}tdd
|d�}|j|jƒ ƒ �qW tdƒ}||d< tdd|d�S t |tƒ�rúg }
x4|D ],\}}td|| j
|ƒd�}|
j|jƒ ƒ �q¬W 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_stringzMysqlx.Datatypes.Anyr	   )rˆ   Úscalarr   rl   z#Mysqlx.Datatypes.Object.ObjectField)Úkeyr†   zMysqlx.Datatypes.Object)Úfld)rˆ   ÚobjzMysqlx.Datatypes.Arrayr†   é   )rˆ   ÚarrayN)Ú
isinstancerM   r   Úboolr   rv   r   r   Útuplerq   Ú_create_anyÚget_messageÚdictÚlistÚitemsrc   )r:   r…   r†   r‰   Zarg_keyÚ	arg_valueÚobj_fldrŒ   Zarray_valuesr–   Úobj_fldsrŠ   Úmsg_objÚmsg_anyrg   r;   r;   r<   r’   p  sl    	




zProtocol._create_anyT)r   Ú	is_scalarr1   c       
         s²   t tttf dœ‡‡fdd„‰ |jƒ }|jƒ }|dkrH‡ fdd„|D ƒS t|ƒ}|dg }|t|ƒkrntdƒ‚x>|jƒ D ]2\}}||kr–td|› �ƒ‚|| }	ˆ |ƒ||	< qxW |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.
        )r†   r1   c                s   ˆ rt | ƒjƒ S ˆj| ƒjƒ S )N)r   r“   r’   )r†   )rœ   r:   r;   r<   Úbuild_valueË  s    z/Protocol._get_binding_args.<locals>.build_valueNc                s   g | ]}ˆ |ƒ‘qS r;   r;   )Ú.0r†   )r�   r;   r<   ú
<listcomp>×  s    z.Protocol._get_binding_args.<locals>.<listcomp>z;The number of bind parameters and placeholders do not matchz*Unable to find placeholder for parameter: )	r   r   r*   r)   Úget_bindingsZget_binding_maprq   r   r–   )
r:   r   rœ   ZbindingsZbinding_mapÚcountÚargsÚnamer†   Úposr;   )r�   rœ   r:   r<   Ú_get_binding_argsº  s$    
zProtocol._get_binding_args)rg   Úresultr1   c             C   s(  |d dkrRt jd|d ƒ}| jj|jƒ tjd|j|jƒ |j|j	|j|jƒ nÒ|d dkrpt jd|d ƒ n´|d dk�r$t jd	|d ƒ}|d
 t
dƒkr¸|jdd„ |d D ƒƒ nlt|d ttƒƒrÖ|d d n|d }|d
 t
dƒk�r|jt|dƒƒ n"|d
 t
dƒk�r$|jt|dƒƒ dS )z¨Process frame.

        Args:
            msg (mysqlx.protobuf.Message): A MySQL X Protobuf Message.
            result (Result): A `Result` based type object.
        rˆ   r	   zMysqlx.Notice.Warningr[   z:Protocol.process_frame Received Warning Notice code %s: %srl   z$Mysqlx.Notice.SessionVariableChangedr�   z!Mysqlx.Notice.SessionStateChangedÚparamzBMysqlx.Notice.SessionStateChanged.Parameter.GENERATED_DOCUMENT_IDSc             S   s    g | ]}t t |d ƒdƒjƒ ‘qS )Zv_octetsr†   )r   Údecode)rž   r†   r;   r;   r<   rŸ     s   z+Protocol._process_frame.<locals>.<listcomp>r†   r   z9Mysqlx.Notice.SessionStateChanged.Parameter.ROWS_AFFECTEDZv_unsigned_intz?Mysqlx.Notice.SessionStateChanged.Parameter.GENERATED_INSERT_IDN)r   Zfrom_messager}   rc   rg   r   ÚwarningÚcodeZappend_warningÚlevelr   Zset_generated_idsr�   r‘   r   Zset_rows_affectedr   Zset_generated_insert_id)r:   rg   r¦   Zwarn_msgZsess_state_msgZsess_state_valuer;   r;   r<   Ú_process_frameè  s:    

zProtocol._process_frame)r¦   r1   c             C   s  �xy| j jƒ }W nF tk
rX } z*t|jƒ ƒ}|rHt|› d|› �ƒ|‚W Y dd}~X nX |jdkrvt|d |d ƒ‚|jdkr®y| j||ƒ W n tt	fk
rª   wY nX q|jdkr¼dS |jdkrÒ|j
d	ƒ q|jd
krè|jd	ƒ q|jdk�r|jd	ƒ P qP qW |S )z`Read message.

        Args:
            result (Result): A `Result` based type object.
        z	 reason: NzMysqlx.Errorrg   rª   zMysqlx.Notice.FramezMysqlx.Sql.StmtExecuteOkzMysqlx.Resultset.FetchDoneTz(Mysqlx.Resultset.FetchDoneMoreResultsetszMysqlx.Resultset.Row)rz   rh   ÚRuntimeErrorÚreprZget_warningsrˆ   r   r¬   ÚAttributeErrorÚKeyErrorZ
set_closedZset_has_more_resultsZset_has_data)r:   r¦   rg   ÚerrÚwarningsr;   r;   r<   rb     s2    &






zProtocol._read_message)r0   r1   c             C   s"   || _ | jj|ƒ | 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)r|   rz   rj   r{   )r:   r0   r;   r;   r<   rj   ?  s    
zProtocol.set_compressionc             C   s^   t dƒ}| jjtdƒ|ƒ | jjƒ }x|jdkr<| jjƒ }q&W |jdkrZt|d |d ƒ‚|S )zkGet capabilities.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        z!Mysqlx.Connection.CapabilitiesGetz/Mysqlx.ClientMessages.Type.CON_CAPABILITIES_GETzMysqlx.Notice.FramezMysqlx.Errorrg   rª   )r   r{   ru   r   rz   rh   rˆ   r   )r:   rg   r;   r;   r<   Úget_capabilitesM  s    

zProtocol.get_capabilites)Úkwargsr1   c             K   s(  |sdS t dƒ}x´|jƒ D ]¨\}}t dƒ}||d< t|tƒr |}g }x2|D ]*}t d|| j|| ƒd�}	|j|	jƒ ƒ qJW t d|d�}
t d	d
|
d�}|jƒ |d< n| j|ƒ|d< |d j|jƒ gƒ qW t dƒ}||d< | jj	t
dƒ|ƒ y| jƒ S  tk
�r" } z|jdk�r‚ W Y dd}~X nX 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£   z#Mysqlx.Datatypes.Object.ObjectField)rŠ   r†   zMysqlx.Datatypes.Object)r‹   zMysqlx.Datatypes.Anyrl   )rˆ   rŒ   r†   Úcapabilitiesz!Mysqlx.Connection.CapabilitiesSetz/Mysqlx.ClientMessages.Type.CON_CAPABILITIES_SETiŠ  )r   r–   r�   r”   r’   rc   r“   rƒ   r{   ru   r   Úread_okr
   Úerrno)r:   r´   rµ   rŠ   r†   Z
capabilityr–   r™   Úitemr˜   rš   r›   rg   r±   r;   r;   r<   Úset_capabilitiesa  s>    	

zProtocol.set_capabilities)ÚmethodÚ	auth_dataÚinitial_responser1   c             C   sF   t dƒ}||d< |dk	r ||d< |dk	r0||d< | j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_nameNr»   r¼   z2Mysqlx.ClientMessages.Type.SESS_AUTHENTICATE_START)r   r{   ru   r   )r:   rº   r»   r¼   rg   r;   r;   r<   Úsend_auth_start‘  s    zProtocol.send_auth_startc             C   s>   | j jƒ }x|jdkr"| j jƒ }qW |jdkr6tdƒ‚|d S )züRead authenticate continue.

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

        Returns:
            str: The authentication data.
        zMysqlx.Notice.Framez#Mysqlx.Session.AuthenticateContinuez>Unexpected message encountered during authentication handshaker»   )rz   rh   rˆ   r
   )r:   rg   r;   r;   r<   Úread_auth_continue©  s    


zProtocol.read_auth_continue)r»   r1   c             C   s"   t d|d�}| jjtdƒ|ƒ dS )zeSend authenticate continue.

        Args:
            auth_data (str): Authentication data.
        z#Mysqlx.Session.AuthenticateContinue)r»   z5Mysqlx.ClientMessages.Type.SESS_AUTHENTICATE_CONTINUEN)r   r{   ru   r   )r:   r»   rg   r;   r;   r<   Úsend_auth_continue¼  s    zProtocol.send_auth_continuec             C   s4   x.| j jƒ }|jdkrP |jdkrt|jƒ‚qW dS )z~Read authenticate OK.

        Raises:
            :class:`mysqlx.InterfaceError`: If message type is `Mysqlx.Error`.
        zMysqlx.Session.AuthenticateOkzMysqlx.ErrorN)rz   rh   rˆ   r
   rg   )r:   rg   r;   r;   r<   Úread_auth_okÈ  s    


zProtocol.read_auth_ok)rf   rg   r   r1   c             C   sR  |j rÂ|jdkrÂ|jdkr*| j|ƒ\}}nB|jdkrD| j|ƒ\}}n(|jdkr^| j|ƒ\}}ntd|› �ƒ‚t|jƒ ƒ}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jtdƒ|ƒ y| jƒ  W n* tk
�rL } zt|‚W Y dd}~X nX 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
        zMysqlx.Crud.InsertzMysqlx.Crud.FindzMysqlx.Crud.UpdatezMysqlx.Crud.DeletezInvalid message type: z!Mysqlx.Expr.Expr.Type.PLACEHOLDERzMysqlx.Crud.LimitExprzMysqlx.Expr.Expr)rˆ   ÚpositionÚ	row_countr	   ÚoffsetZ
limit_exprz#Mysqlx.Prepare.Prepare.OneOfMessagerˆ   zMysqlx.Prepare.PrepareÚstmt_idr   z*Mysqlx.ClientMessages.Type.PREPARE_PREPAREN)Ú	has_limitrˆ   Ú
build_findÚbuild_updateÚbuild_deletera   rq   r    r   r   r   rÄ   r{   ru   r¶   r
   r   )r:   rf   rg   r   Ú_rÁ   ÚplaceholderZmsg_limit_exprÚ
oneof_typeÚoneof_opÚ	msg_oneofZmsg_preparer±   r;   r;   r<   Úsend_prepare_prepareÕ  s>    




zProtocol.send_prepare_preparec       	      C   s¤   t | \}}tdƒ}t|ƒ|d< |||< tdƒ}|j|d< | j|dd�}|rZ|d j|ƒ |jrŽ|d j| j|jƒ ƒj	ƒ | j|j
ƒ ƒj	ƒ gƒ | j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
        z#Mysqlx.Prepare.Prepare.OneOfMessagerˆ   zMysqlx.Prepare.ExecuterÄ   F)rœ   r¢   z*Mysqlx.ClientMessages.Type.PREPARE_EXECUTEN)r   r   r   rÄ   r¥   rƒ   rÅ   r’   Úget_limit_row_countr“   Úget_limit_offsetr{   ru   )	r:   rf   rg   r   rË   rÌ   rÍ   Zmsg_executer¢   r;   r;   r<   Úsend_prepare_execute  s     
zProtocol.send_prepare_execute)rÄ   r1   c             C   s.   t dƒ}||d< | jjtdƒ|ƒ | jƒ  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   r{   ru   r   r¶   )r:   rÄ   Zmsg_deallocr;   r;   r<   Úsend_prepare_deallocate>  s    	z Protocol.send_prepare_deallocatec             C   sp   |j r8tdƒ}|jƒ |d< |jdkr0|jƒ |d< ||d< |dk}| j||d�}|r`|d j|ƒ | j||ƒ 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Â   zMysqlx.Crud.FindrÃ   Úlimitz+Mysqlx.ClientMessages.Type.SQL_STMT_EXECUTE)rœ   r¢   N)rÅ   r   rÏ   rˆ   rÐ   r¥   rƒ   Úsend_msg)r:   rf   rg   r   Z	msg_limitrœ   r¢   r;   r;   r<   Úsend_msg_without_psO  s    
zProtocol.send_msg_without_ps)rf   rg   r1   c             C   s   | j 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)r{   ru   r   )r:   rf   rg   r;   r;   r<   rÔ   k  s    
zProtocol.send_msg)r   r1   c             C   s    t |jƒ rdndƒ}td|jj|jjd�}td||d�}|jrJ|jƒ |d< | j||ƒ |j	ƒ rlt dƒ|d	< n|j
ƒ r€t d
ƒ|d	< |jjdkr˜|j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
        zMysqlx.Crud.DataModel.DOCUMENTzMysqlx.Crud.DataModel.TABLEzMysqlx.Crud.Collection)r£   ÚschemazMysqlx.Crud.Find)Ú
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_exprr„   Zis_lock_exclusiveZis_lock_sharedZlock_contentionr†   )r:   r   r×   rØ   rg   r;   r;   r<   rÆ   w  s$    zProtocol.build_findc             C   s®   t |jƒ rdndƒ}td|jj|jjd�}td||d�}| j||ƒ x`|jƒ jƒ D ]P\}}tdƒ}|j	|d< |j
|d	< |jd
k	rŽt|jƒ|d< |d j|jƒ gƒ qRW d|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
        zMysqlx.Crud.DataModel.DOCUMENTzMysqlx.Crud.DataModel.TABLEzMysqlx.Crud.Collection)r£   rÖ   zMysqlx.Crud.Update)r×   rØ   zMysqlx.Crud.UpdateOperationÚ	operationÚsourceNr†   z&Mysqlx.ClientMessages.Type.CRUD_UPDATE)r   rÚ   r   rÛ   r£   rÖ   r„   Zget_update_opsr–   Zupdate_typerÝ   r†   r   rƒ   r“   )r:   r   r×   rØ   rg   rÉ   Z	update_oprÜ   r;   r;   r<   rÇ   ¡  s$    


zProtocol.build_updatec             C   sL   t |jƒ rdndƒ}td|jj|jjd�}td||d�}| j||ƒ 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
        zMysqlx.Crud.DataModel.DOCUMENTzMysqlx.Crud.DataModel.TABLEzMysqlx.Crud.Collection)r£   rÖ   zMysqlx.Crud.Delete)r×   rØ   z&Mysqlx.ClientMessages.Type.CRUD_DELETE)r   rÚ   r   rÛ   r£   rÖ   r„   )r:   r   r×   rØ   rg   r;   r;   r<   rÈ   Ê  s    zProtocol.build_delete)Ú	namespacer   Úfieldsr1   c             C   s€   t d||dd�}|rxg }x6|jƒ D ]*\}}t d|| j|ƒd�}|j|jƒ ƒ q"W t d|d�}	t dd	|	d
�}
|
jƒ 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)rÞ   r   Zcompact_metadataz#Mysqlx.Datatypes.Object.ObjectField)rŠ   r†   zMysqlx.Datatypes.Object)r‹   zMysqlx.Datatypes.Anyrl   )rˆ   rŒ   r¢   z+Mysqlx.ClientMessages.Type.SQL_STMT_EXECUTE)r   r–   r’   rc   r“   )r:   rÞ   r   rß   rg   r™   rŠ   r†   r˜   rš   r›   r;   r;   r<   Úbuild_execute_statementë  s"    z Protocol.build_execute_statementc       	      C   s  t | jƒ rdndƒ}td| jj| jjd�}td||d�}t| dƒrzx6| jD ],}t|| jƒ  ƒj	ƒ }|d j
|jƒ gƒ qJW xv| jƒ D ]j}td	ƒ}t|tƒrÂx>|D ]}|d
 j
t|ƒjƒ gƒ q W n|d
 j
t|ƒjƒ gƒ |d j
|jƒ gƒ q„W t| dƒ�r
| jƒ |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
        zMysqlx.Crud.DataModel.DOCUMENTzMysqlx.Crud.DataModel.TABLEzMysqlx.Crud.Collection)r£   rÖ   zMysqlx.Crud.Insert)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_fieldrƒ   r“   Z
get_valuesr�   r•   r   rä   )	r   r×   rØ   rg   râ   Úexprr†   rã   Úvalr;   r;   r<   Úbuild_insert  s0    


zProtocol.build_insertc             C   s   | j |ƒ}|dk	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)rb   r   )r:   r¦   rg   r;   r;   r<   Úclose_resultJ  s    	
zProtocol.close_resultc             C   s4   | j |ƒ}|dkrdS |jdkr$|S | jj|ƒ dS )z\Read row.

        Args:
            result (Result): A `Result` based type object.
        NzMysqlx.Resultset.Row)rb   rˆ   rz   ri   )r:   r¦   rg   r;   r;   r<   Úread_rowW  s    

zProtocol.read_rowc             C   s¶   g }x¬| j |ƒ}|dkrP |jdkr2| jj|ƒ P |jdkrDtdƒ‚t|d |d |d |d |d	 |d
 |d |jddƒ|jddƒ|jddƒ|jddƒ|jdƒƒ}|j|ƒ qW |S )z¿Returns column metadata.

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

        Raises:
            :class:`mysqlx.InterfaceError`: If unexpected message.
        NzMysqlx.Resultset.RowzMysqlx.Resultset.ColumnMetaDatazUnexpected msg typerˆ   ÚcatalogrÖ   ÚtableZoriginal_tabler£   Úoriginal_nameÚlengthé   Z	collationr   Zfractional_digitsÚflagsé   Úcontent_type)rb   rˆ   rz   ri   r
   r   Úgetrc   )r:   r¦   Úcolumnsrg   Úcolr;   r;   r<   Úget_column_metadatae  s2    	






zProtocol.get_column_metadatac             C   sD   | j 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.
        zMysqlx.ErrorzMysqlx.Error: rg   rª   )r·   z	Mysqlx.OkzUnexpected message encounteredN)rz   rh   rˆ   r
   )r:   rg   r;   r;   r<   r¶   ‰  s
    


zProtocol.read_okc             C   s   t dƒ}| jjtdƒ|ƒ dS )zSend connection close.zMysqlx.Connection.Closez$Mysqlx.ClientMessages.Type.CON_CLOSEN)r   r{   ru   r   )r:   rg   r;   r;   r<   Úsend_connection_close•  s    zProtocol.send_connection_closec             C   s   t dƒ}| jjtdƒ|ƒ dS )zSend close.zMysqlx.Session.Closez%Mysqlx.ClientMessages.Type.SESS_CLOSEN)r   r{   ru   r   )r:   rg   r;   r;   r<   Ú
send_closeœ  s    zProtocol.send_closec             C   sL   t dƒ}tdƒ}||d< d|d< tdƒ}|jƒ g|d< | j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“   r{   ru   )r:   Zcond_keyZmsg_ocZmsg_eor;   r;   r<   Úsend_expect_open£  s    zProtocol.send_expect_open)Ú	keep_openr1   c             C   st   t dƒ}|dkrBy| jƒ  | jƒ  d}W n tk
r@   d}Y nX |rNd|d< | jjtdƒ|ƒ | jƒ  |rpdS dS )z¨Send reset session message.

        Returns:
            boolean: ``True`` if the server will keep the session open,
                     otherwise ``False``.
        zMysqlx.Session.ResetNTFrú   z%Mysqlx.ClientMessages.Type.SESS_RESET)r   rù   r¶   r
   r{   ru   r   )r:   rú   rg   r;   r;   r<   Ú
send_reset±  s     
zProtocol.send_reset)T)NN)N)N)BrI   rJ   rK   rL   rO   rk   r=   Úpropertyr   rM   r~   Ústaticmethodr(   r   r„   r   r’   r   r%   r�   r   r*   r)   r¥   r+   r¬   rb   rj   r³   r¹   r½   rN   r¾   r¿   rÀ   r    r   r"   r#   r$   r&   rÎ   rÑ   rv   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<   rw   J  sr   K&,6!3 9#

)
(
#%
2$rw   )DrL   r^   r7   Úior   Útypingr   r   r   r   r   r   Z	lz4.framerA   ZHAVE_LZ4ÚImportErrorZ	zstandardr6   Z	HAVE_ZSTDÚerrorsr
   r   r   r   ræ   r   r   r   r   r   r   Zhelpersr   r   r   Zprotobufr   r   r   r   r   r¦   r   Z	statementr   r   r   r    r!   r"   r#   r$   r%   r&   Útypesr'   r(   r)   r*   r+   r,   r-   r.   rn   r/   rO   rk   rw   r;   r;   r;   r<   Ú<module>   s6    

 0(Dd=