o
    –¨Êh3a  ã                   @   s¨   d Z ddlZddlmZmZ ddl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 G d
d„ deƒZG dd„ deƒZdS )z3Implementation of the X protocol for MySQL servers.é    Né   )ÚSTRING_TYPESÚ	INT_TYPES)ÚInterfaceErrorÚ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)ÚSERVER_MESSAGESÚPROTOBUF_REPEATED_TYPESÚMessageÚmysqlxpb_enumc                   @   s8   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
d„ ZdS )ÚMessageReaderWriterz‚Implements a Message Reader/Writer.

    Args:
        socket_stream (mysqlx.connection.SocketStream): `SocketStream` object.
    c                 C   s   || _ d | _d S ©N)Ú_streamÚ_msg)ÚselfÚsocket_stream© r   úA/var/www/html/env/lib/python3.10/site-packages/mysqlx/protocol.pyÚ__init__1   s   
zMessageReaderWriter.__init__c                 C   sž   | j  d¡}t d|¡\}}|dkrtdƒ‚| j  |d ¡}t |¡}|s,td |¡ƒ‚|dkr8|dkr8|  	¡ S z	t
 ||¡}W |S  tyN   |  	¡  Y S w )	aK  Read message.

        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.
        é   ú<LBé
   zZThe connected server does not have the MySQL X protocol plugin enabled or protocolmismatchr   zUnknown msg_type: {0}é   ó    )r   ÚreadÚstructÚunpackr   r   ÚgetÚ
ValueErrorÚformatÚ_read_messager   Úfrom_server_messageÚRuntimeError)r   ÚhdrÚmsg_lenÚmsg_typeÚpayloadÚmsg_type_nameÚmsgr   r   r   r)   5   s    
þÿz!MessageReaderWriter._read_messagec                 C   s"   | j dur| j }d| _ |S |  ¡ S )zgRead message.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        N)r   r)   ©r   r1   r   r   r   Úread_messageT   s
   
z MessageReaderWriter.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)r   r   r2   r   r   r   Úpush_message`   s   
	
z MessageReaderWriter.push_messagec                 C   s<   t | ¡ ƒ}t dt|ƒd |¡}| j d ||g¡¡ dS )z•Write message.

        Args:
            msg_id (int): The message ID.
            msg (mysqlx.protobuf.Message): MySQL X Protobuf Message.
        r   r   r"   N)r   Úserialize_to_stringr$   ÚpackÚlenr   ÚsendallÚjoin)r   Úmsg_idr1   Úmsg_strÚheaderr   r   r   Úwrite_messagem   s   z!MessageReaderWriter.write_messageN)	Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r)   r3   r4   r=   r   r   r   r   r   +   s    r   c                   @   sÒ   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
d„ Zdd„ Z	dd„ Z
d3d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d)d*„ Zd+d,„ Zd-d.„ Zd/d0„ Zd1d2„ ZdS )4ÚProtocolz—Implements the MySQL X Protocol.

    Args:
        read_writer (mysqlx.protocol.MessageReaderWriter): A Message             Reader/Writer object.
    c                 C   s   || _ || _d | _d S r   )Ú_readerÚ_writerÚ_message)r   Úreader_writerr   r   r   r   €   s   
zProtocol.__init__c                 C   s–   |j r	| ¡ |d< |jr|d  |  |¡¡ |jr&td| ¡ | ¡ d�|d< |j	r2|d  | 
¡ ¡ |jr>|d  | ¡ ¡ |jrI| ¡ |d< d	S d	S )
z­Apply filter.

        Args:
            msg (mysqlx.protobuf.Message): The MySQL X Protobuf Message.
            stmt (Statement): A `Statement` based type object.
        ÚcriteriaÚargszMysqlx.Crud.Limit)Ú	row_countÚoffsetÚlimitÚorderÚgroupingÚgrouping_criteriaN)Ú	has_whereÚget_where_exprÚhas_bindingsÚextendÚget_binding_scalarsÚ	has_limitr   Úget_limit_row_countÚget_limit_offsetÚhas_sortÚget_sort_exprÚhas_group_byÚget_groupingÚ
has_havingÚ
get_having)r   r1   Ústmtr   r   r   Ú_apply_filter…   s"   
ýÿzProtocol._apply_filterc                 C   sŒ  t |tƒrtd|d�}tdd|d�}tdd|d�S t |tƒr'tddt|ƒd�S t |tƒrB|d	k r9tddt|ƒd�S tddt|ƒd�S t |tƒrkt	|ƒd
krk|\}}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 ]2}	g }
|	 ¡ D ]\}}td||  
|¡d�}|
 | ¡ ¡ qŠtd|
d�}td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é   )ÚtypeÚv_stringúMysqlx.Datatypes.Anyr   )ra   Úscalarr   é   ú#Mysqlx.Datatypes.Object.ObjectField©Úkeyr_   úMysqlx.Datatypes.Object©Úfld©ra   ÚobjzMysqlx.Datatypes.Arrayr_   é   )ra   ÚarrayN)Ú
isinstancer   r   Úboolr   r   r   r   Útupler7   Ú_create_anyÚget_messageÚdictÚlistÚitemsÚappend)r   Úargr_   rd   Úarg_keyÚ	arg_valueÚobj_fldrm   Úarray_valuesrw   Úobj_fldsrh   Úmsg_objÚmsg_anyr1   r   r   r   rs   œ   sV   
	
ÿ
ÿÿ
ÿÿÿ
ÿzProtocol._create_anyc                 C   s  |d dkrt  d|d ¡}| |j|j|j¡ dS |d dkr*t  d|d ¡ dS |d dkr…t  d|d ¡}|d	 td
ƒkrN| dd„ |d D ƒ¡ dS t|d t	t
ƒƒr]|d d n|d }|d	 tdƒkrs| t|dƒ¡ dS |d	 tdƒkr‡| t|dƒ¡ dS dS dS )z¨Process frame.

        Args:
            msg (mysqlx.protobuf.Message): A MySQL X Protobuf Message.
            result (Result): A `Result` based type object.
        ra   r   zMysqlx.Notice.Warningr/   re   z$Mysqlx.Notice.SessionVariableChangedrn   z!Mysqlx.Notice.SessionStateChangedÚparamzBMysqlx.Notice.SessionStateChanged.Parameter.GENERATED_DOCUMENT_IDSc                 S   s    g | ]}t t |d ƒdƒ ¡ ‘qS )Úv_octetsr_   )r   Údecode)Ú.0r_   r   r   r   Ú
<listcomp>ä   s    þ
ÿÿz+Protocol._process_frame.<locals>.<listcomp>r_   r   z9Mysqlx.Notice.SessionStateChanged.Parameter.ROWS_AFFECTEDÚv_unsigned_intz?Mysqlx.Notice.SessionStateChanged.Parameter.GENERATED_INSERT_IDN)r   Úfrom_messageÚappend_warningÚlevelÚcoder1   r   Úset_generated_idsrp   rr   r   Úset_rows_affectedr   Úset_generated_insert_id)r   r1   ÚresultÚwarn_msgÚsess_state_msgÚsess_state_valuer   r   r   Ú_process_frameÏ   sR   ÿÿÿÿþÿÿÿýÿÿÿ
ÿézProtocol._process_framec                 C   sª   	 | j  ¡ }|jdkrt|d |d ƒ‚|jdkr'z|  ||¡ W n2   Y q |jdkr.dS |jdkr9| d¡ n|jd	krD| d¡ n|jd
krQ| d¡ 	 |S 	 |S q)z`Read message.

        Args:
            result (Result): A `Result` based type object.
        TúMysqlx.Errorr1   rŠ   úMysqlx.Notice.FramezMysqlx.Sql.StmtExecuteOkNzMysqlx.Resultset.FetchDonez(Mysqlx.Resultset.FetchDoneMoreResultsetsúMysqlx.Resultset.Row)rC   r3   ra   r   r’   Ú
set_closedÚset_has_more_resultsÚset_has_data©r   rŽ   r1   r   r   r   r)   ÷   s,   







ÿìzProtocol._read_messagec                 C   sF   t dƒ}| j tdƒ|¡ | j ¡ }|jdkr!| j ¡ }|jdks|S )zkGet capabilities.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        z!Mysqlx.Connection.CapabilitiesGetz/Mysqlx.ClientMessages.Type.CON_CAPABILITIES_GETr”   )r   rD   r=   r   rC   r3   ra   r2   r   r   r   Úget_capabilites  s   þ



ÿzProtocol.get_capabilitesc                 K   sv   t dƒ}| ¡ D ]\}}t dƒ}||d< |  |¡|d< |d  | ¡ g¡ qt dƒ}||d< | j tdƒ|¡ |  ¡ S )z­Set capabilities.

        Args:
            **kwargs: Arbitrary keyword arguments.

        Returns:
            mysqlx.protobuf.Message: MySQL X Protobuf Message.
        zMysqlx.Connection.CapabilitieszMysqlx.Connection.CapabilityÚnamer_   Úcapabilitiesz!Mysqlx.Connection.CapabilitiesSetz/Mysqlx.ClientMessages.Type.CON_CAPABILITIES_SET)	r   rw   rs   rR   rt   rD   r=   r   Úread_ok)r   Úkwargsrœ   rh   r_   Ú
capabilityr1   r   r   r   Úset_capabilities"  s   	þzProtocol.set_capabilitiesNc                 C   sF   t dƒ}||d< |dur||d< |dur||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.AuthenticateStartÚ	mech_nameNÚ	auth_dataÚinitial_responsez2Mysqlx.ClientMessages.Type.SESS_AUTHENTICATE_START©r   rD   r=   r   )r   Úmethodr¢   r£   r1   r   r   r   Úsend_auth_start8  s   ÿÿzProtocol.send_auth_startc                 C   sB   | j  ¡ }|jdkr| j  ¡ }|jdks
|jdkrt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¢   )rC   r3   ra   r   r2   r   r   r   Úread_auth_continueI  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¢   r1   r   r   r   Úsend_auth_continue[  s   ÿÿÿzProtocol.send_auth_continuec                 C   s0   	 | j  ¡ }|jdkrdS |jdkrt|jƒ‚q)z~Read authenticate OK.

        Raises:
            :class:`mysqlx.InterfaceError`: If message type is `Mysqlx.Error`.
        TzMysqlx.Session.AuthenticateOkr“   N)rC   r3   ra   r   r1   r2   r   r   r   Úread_auth_okf  s   



ûzProtocol.read_auth_okc           	      C   s~   |  ¡ }| ¡ }t|ƒ}|dg }|t|ƒkrtdƒ‚|D ]}|d }||vr.td |¡ƒ‚|| }t|d ƒ ¡ ||< q|S )a  Returns the binding scalars.

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

        Returns:
            list: A list of ``mysqlx.protobuf.Message`` objects.
        Nz;The number of bind parameters and placeholders do not matchr›   z-Unable to find placeholder for parameter: {0}r_   )Úget_bindingsÚget_binding_mapr7   r   r(   r
   rt   )	r   r]   ÚbindingsÚbinding_mapÚcountÚscalarsÚbindingr›   Úposr   r   r   rS   s  s   

ÿzProtocol.get_binding_scalarsc                 C   sª   t | ¡ rdndƒ}td|jj|jjd�}td||d�}|jr%| ¡ |d< |  ||¡ | 	¡ r6t dƒ|d	< n
| 
¡ r@t d
ƒ|d	< |jdkrJ|j|d< | j t dƒ|¡ dS )z§Send find.

        Args:
            stmt (Statement): A :class:`mysqlx.ReadStatement` or
                              :class:`mysqlx.FindStatement` object.
        úMysqlx.Crud.DataModel.DOCUMENTúMysqlx.Crud.DataModel.TABLEúMysqlx.Crud.Collection©r›   ÚschemazMysqlx.Crud.Find©Ú
data_modelÚ
collectionÚ
projectionz'Mysqlx.Crud.Find.RowLock.EXCLUSIVE_LOCKÚlockingz$Mysqlx.Crud.Find.RowLock.SHARED_LOCKr   Úlocking_optionsz$Mysqlx.ClientMessages.Type.CRUD_FINDN)r   Úis_doc_basedr   Útargetr›   r·   Úhas_projectionÚget_projection_exprr^   Úis_lock_exclusiveÚis_lock_sharedÚlock_contentionrD   r=   ©r   r]   r¹   rº   r1   r   r   r   Ú	send_find�  s4   ÿþþÿÿÿ

ÿzProtocol.send_findc                 C   s°   t | ¡ rdndƒ}td|jj|jjd�}td||d�}|  ||¡ | ¡ D ]&}tdƒ}|j|d< |j	|d	< |j
d
urBt|j
ƒ|d< |d  | ¡ g¡ q&| j t dƒ|¡ d
S )z­Send update.

        Args:
            stmt (Statement): A :class:`mysqlx.ModifyStatement` or
                              :class:`mysqlx.UpdateStatement` object.
        r³   r´   rµ   r¶   zMysqlx.Crud.Updater¸   zMysqlx.Crud.UpdateOperationÚ	operationÚsourceNr_   z&Mysqlx.ClientMessages.Type.CRUD_UPDATE)r   r¾   r   r¿   r›   r·   r^   Úget_update_opsÚupdate_typerÈ   r_   r	   rR   rt   rD   r=   )r   r]   r¹   rº   r1   Ú	update_oprÇ   r   r   r   Úsend_update­  s.   ÿþþÿ


ÿzProtocol.send_updatec                 C   sZ   t | ¡ rdndƒ}td|jj|jjd�}td||d�}|  ||¡ | j t dƒ|¡ dS )	z­Send delete.

        Args:
            stmt (Statement): A :class:`mysqlx.DeleteStatement` or
                              :class:`mysqlx.RemoveStatement` object.
        r³   r´   rµ   r¶   zMysqlx.Crud.Deleter¸   z&Mysqlx.ClientMessages.Type.CRUD_DELETEN)	r   r¾   r   r¿   r›   r·   r^   rD   r=   rÅ   r   r   r   Úsend_deleteÈ  s   ÿþ
ÿÿÿzProtocol.send_deletec                 C   sÖ   t d||dd�}|dkrLt|ttfƒr|d  ¡ n| ¡ }g }|D ]\}}t d||  |¡d�}	| |	 ¡ ¡ q!t d|d	�}
t d
d|
d�}| ¡ g|d< n|D ]}|  |¡}|d  | ¡ g¡ qN| j	 
tdƒ|¡ dS )zËSend execute statement.

        Args:
            namespace (str): The namespace.
            stmt (Statement): A `Statement` based type object.
            args (iterable): An iterable object.
        zMysqlx.Sql.StmtExecuteF)Ú	namespacer]   Úcompact_metadataÚmysqlxr   rf   rg   ri   rj   rc   re   rl   rH   z+Mysqlx.ClientMessages.Type.SQL_STMT_EXECUTEN)r   rp   rv   rr   rw   rs   rx   rt   rR   rD   r=   r   )r   rÎ   r]   rH   r1   rw   r~   rh   r_   r|   r   r€   ry   r   r   r   Úsend_execute_statementÚ  s,   ÿÿ
ÿ
ÿzProtocol.send_execute_statementc           
      C   s  t | ¡ rdndƒ}td|jj|jjd�}td||d�}t|dƒr;|jD ]}t|| ¡  ƒ 	¡ }|d  
| ¡ g¡ q$| ¡ D ]3}td	ƒ}t|tƒr\|D ]}	|d
  
t|	ƒ ¡ g¡ qLn|d
  
t|ƒ ¡ g¡ |d  
| ¡ g¡ q?t|dƒr~| ¡ |d< | j t dƒ|¡ dS )zªSend insert.

        Args:
            stmt (Statement): A :class:`mysqlx.AddStatement` or
                              :class:`mysqlx.InsertStatement` object.
        r³   r´   rµ   r¶   zMysqlx.Crud.Insertr¸   Ú_fieldsr»   zMysqlx.Crud.Insert.TypedRowÚfieldÚrowÚ	is_upsertÚupsertz&Mysqlx.ClientMessages.Type.CRUD_INSERTN)r   r¾   r   r¿   r›   r·   ÚhasattrrÒ   r   Úparse_table_insert_fieldrR   rt   Ú
get_valuesrp   rv   r	   rÕ   rD   r=   )
r   r]   r¹   rº   r1   rÓ   Úexprr_   rÔ   Úvalr   r   r   Úsend_insertú  s>   ÿþþÿ

ÿ
ÿ
ÿzProtocol.send_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)r)   r   r™   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•   )r)   ra   rC   r4   r™   r   r   r   Úread_row+  s   

zProtocol.read_rowc                 C   s¸   g }	 |   |¡}|du r	 |S |jdkr| j |¡ 	 |S |jdkr&tdƒ‚t|d |d |d |d	 |d
 |d |d | dd¡| dd¡| dd¡| dd¡| d¡ƒ}| |¡ q)z¿Returns column metadata.

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

        Raises:
            :class:`mysqlx.InterfaceError`: If unexpected message.
        TNr•   zMysqlx.Resultset.ColumnMetaDatazUnexpected msg typera   Úcatalogr·   ÚtableÚoriginal_tabler›   Úoriginal_nameÚlengthé   Ú	collationr   Úfractional_digitsÚflagsé   Úcontent_type)r)   ra   rC   r4   r   r   r&   rx   )r   rŽ   Úcolumnsr1   Úcolr   r   r   Úget_column_metadata9  s.   	

ò
õ



ù
ïzProtocol.get_column_metadatac                 C   sF   | j  ¡ }|jdkrtd |d ¡ƒ‚|jdkr!td |d ¡ƒ‚dS )zeRead OK.

        Raises:
            :class:`mysqlx.InterfaceError`: If unexpected message.
        r“   zMysqlx.Error: {}r1   z	Mysqlx.Okz"Unexpected message encountered: {}N)rC   r3   ra   r   r(   r2   r   r   r   r�   W  s   



ÿÿzProtocol.read_okc                 C   ó   t dƒ}| j tdƒ|¡ dS )zSend connection close.zMysqlx.Connection.Closez$Mysqlx.ClientMessages.Type.CON_CLOSENr¤   r2   r   r   r   Úsend_connection_closed  ó   ÿÿzProtocol.send_connection_closec                 C   rí   )zSend close.zMysqlx.Session.Closez%Mysqlx.ClientMessages.Type.SESS_CLOSENr¤   r2   r   r   r   Ú
send_closej  rï   zProtocol.send_closec                 C   rí   )zSend reset.zMysqlx.Session.Resetz%Mysqlx.ClientMessages.Type.SESS_RESETNr¤   r2   r   r   r   Ú
send_resetp  rï   zProtocol.send_reset)NN)r>   r?   r@   rA   r   r^   rs   r’   r)   rš   r    r¦   r¨   r©   rª   rS   rÆ   rÌ   rÍ   rÑ   rÜ   rÝ   rÞ   rì   r�   rî   rð   rñ   r   r   r   r   rB   y   s4    3(
  $rB   )rA   r$   Úcompatr   r   Úerrorsr   r   r   rÚ   r   r	   r
   r   r   r   Úhelpersr   r   rŽ   r   Úprotobufr   r   r   r   Úobjectr   rB   r   r   r   r   Ú<module>   s    N