3
]ð]<0  ã               @   sj  d Z ddlZddlZddlZddlZddlZdZdZy(ddlmZ ej	ej
B ejB ejB ZW n ek
rt   dZY nX yddlmZ W n ek
rž   eZY nX ddlmZ ddlmZ dd	lmZmZ dd
lmZ ddlmZmZ ddlmZmZm Z m!Z! ddl"m#Z# ej$dƒj%Z&ddd„Z'ej$dƒj%Z(efdd„Z)e�sFdd„ Z*ndd„ Z*dd„ Z+G dd„ de,ƒZ-dS )z&Internal network layer helper methods.é    NT)ÚpollF)Úerror)Ú_decode_all_selective)ÚPY3)ÚhelpersÚmessage)ÚMAX_MESSAGE_SIZE)Ú
decompressÚ_NO_COMPRESSION)ÚAutoReconnectÚNotMasterErrorÚOperationFailureÚProtocolError)Ú_UNPACK_REPLYz<iiiic       (      C   s  t t|ƒƒ}|d }|rdnd}|}|r:| r:tj||ƒ}|r‚|oF|j r‚|jrZ|j|d< |r‚|jjr‚|j	dk	r‚|j	|j
di ƒd< |dk	r’||d< |dk	ož|j}|r®tjjƒ }|rÂ|jƒ tkrÂd}|rð|jrð|jj rð|jj||||ƒ }}d}|�rN|rþd	nd}tj||||||||d
�\}}}}|�rn|dk	�rn||k�rntj|||ƒ n tj||dd|d|||ƒ	\}}}|dk	�rœ||tj k�rœtj|||tj ƒ |�rÊtjjƒ | } |j||||ƒ tjjƒ }yr| j|ƒ |�rð|�rðd}!ddi}"nJt| |ƒ}!|!j||d�}#|#d }"|�r"|j|"|ƒ |	�r:tj|"d|
|d� W nr tk
�r® }$ zT|�rœtjjƒ | |  }%t|$t t!fƒ�r€|$j"}&n
tj#|$ƒ}&|j$|%|&|||ƒ ‚ W Y dd}$~$X nX |�rÚtjjƒ | |  }%|j%|%|"|||ƒ |�r|j�r|!�r|jj&|!j'ƒ ƒ}'t(|'||ƒd }"|"S )a  Execute a command over the socket, or raise socket.error.

    :Parameters:
      - `sock`: a raw socket instance
      - `dbname`: name of the database on which to run the command
      - `spec`: a command document as an ordered dict type, eg SON.
      - `slave_ok`: whether to set the SlaveOkay wire protocol bit
      - `is_mongos`: are we connected to a mongos?
      - `read_preference`: a read preference
      - `codec_options`: a CodecOptions instance
      - `session`: optional ClientSession instance.
      - `client`: optional MongoClient instance for updating $clusterTime.
      - `check`: raise OperationFailure if there are errors
      - `allowable_errors`: errors to ignore if `check` is True
      - `address`: the (host, port) of `sock`
      - `check_keys`: if True, check `spec` for invalid keys
      - `listeners`: An instance of :class:`~pymongo.monitoring.EventListeners`
      - `max_bson_size`: The maximum encoded bson size for this server
      - `read_concern`: The read concern for this command.
      - `parse_write_concern_error`: Whether to parse the ``writeConcernError``
        field in the command response.
      - `collation`: The collation for this command.
      - `compression_ctx`: optional compression Context.
      - `use_op_msg`: True if we should use OP_MSG.
      - `unacknowledged`: True if this is an unacknowledged command.
      - `user_fields` (optional): Response fields that should be decoded
        using the TypeDecoders from codec_options, passed to
        bson._decode_all_selective.
    z.$cmdé   r   ZreadConcernNZafterClusterTimeÚ	collationFé   )Úctxé   Úok)Úcodec_optionsÚuser_fields)Úparse_write_concern_erroréÿÿÿÿ))ÚnextÚiterr   Z_maybe_add_read_preferenceZin_transactionÚlevelÚdocumentÚoptionsZcausal_consistencyZoperation_timeÚ
setdefaultZenabled_for_commandsÚdatetimeÚnowÚlowerr
   Z
_encrypterZ_bypass_auto_encryptionZencryptZ_op_msgZ_raise_document_too_largeÚqueryZ_COMMAND_OVERHEADZpublish_command_startÚsendallÚreceive_messageZunpack_responseZ_process_responser   Z_check_command_responseÚ	ExceptionÚ
isinstancer   r   ÚdetailsZ_convert_exceptionZpublish_command_failureZpublish_command_successZdecryptZraw_command_responser   )(ÚsockZdbnameÚspecZslave_okZ	is_mongosZread_preferencer   ÚsessionÚclientÚcheckZallowable_errorsÚaddressZ
check_keysZ	listenersZmax_bson_sizeZread_concernr   r   Zcompression_ctxZ
use_op_msgZunacknowledgedr   ÚnameÚnsÚflagsÚorigÚpublishÚstartÚ
request_idÚmsgÚsizeZmax_doc_sizeZencoding_durationZreplyZresponse_docZunpacked_docsÚexcÚdurationZfailureZ	decrypted© r:   ú2/tmp/pip-build-20mum3z4/pymongo/pymongo/network.pyÚcommand5   s˜    (














r<   z<iiBc       
      C   sâ   t t| dƒƒ\}}}}|dk	r6||kr6td||f ƒ‚|dkrLtd|f ƒ‚||krdtd||f ƒ‚|dkr–tt| dƒƒ\}}}tt| |d ƒ|ƒ}nt| |d ƒ}yt| }	W n( tk
rØ   td	|tjƒ f ƒ‚Y nX |	|ƒS )
z1Receive a raw BSON message or raise socket.error.é   Nz"Got response id %r but expected %rzEMessage length (%r) not longer than standard message header size (16)z?Message length (%r) is larger than server max message size (%r)iÜ  é	   é   zGot opcode %r but expected %r)Ú_UNPACK_HEADERÚ_receive_data_on_socketr   Ú_UNPACK_COMPRESSION_HEADERr	   r   ÚKeyErrorÚkeys)
r)   r5   Zmax_message_sizeÚlengthÚ_Zresponse_toZop_codeZcompressor_idÚdataZunpack_replyr:   r:   r;   r%   À   s0    
r%   c             C   s¢   t |ƒ}d}xŒ|r˜y| j|ƒ}W n8 ttfk
rX } zt|ƒtjkrFw‚ W Y d d }~X nX |dkrjtdƒ‚||||t|ƒ …< |t|ƒ7 }|t|ƒ8 }qW t	|ƒS )Nr   ó    zconnection closed)
Ú	bytearrayÚrecvÚIOErrorÚOSErrorÚ_errno_from_exceptionÚerrnoÚEINTRr   ÚlenÚbytes)r)   rE   ÚbufÚiÚchunkr8   r:   r:   r;   rA   æ   s    rA   c             C   sŽ   t |ƒ}t|ƒ}d}xt||k rˆy| j||d … ƒ}W n8 ttfk
rl } zt|ƒtjkrZw‚ W Y d d }~X nX |dkr~tdƒ‚||7 }qW |S )Nr   zconnection closed)	rI   Ú
memoryviewÚ	recv_intorK   rL   rM   rN   rO   r   )r)   rE   rR   ÚmvÚ
bytes_readZchunk_lengthr8   r:   r:   r;   rA   ù   s    
c             C   s(   t | dƒr| jS | jr | jd S d S d S )NrN   r   )ÚhasattrrN   Úargs)r8   r:   r:   r;   rM     s
    

rM   c               @   s   e Zd Zdd„ Zdd„ ZdS )ÚSocketCheckerc             C   s(   t rtjƒ | _tƒ | _nd | _d | _d S )N)Ú	_HAS_POLLÚ	threadingÚLockÚ_lockr   Ú_poller)Úselfr:   r:   r;   Ú__init__  s
    

zSocketChecker.__init__c             C   sî   xèyd| j rL| j�4 | j j|tƒ z| j jdƒ}W d| j j|ƒ X W dQ R X ntj|gg g dƒ\}}}W nv ttfk
r€   ‚ Y n^ t	k
r’   dS  t
tfk
rÊ } zt|ƒtjtjfkr¼wdS d}~X n tk
rÜ   dS X t|ƒdkS dS )zHReturn True if we know socket has been closed, False otherwise.
        r   NT)r`   r_   ÚregisterÚ_EVENT_MASKr   Ú
unregisterÚselectÚRuntimeErrorrC   Ú
ValueErrorÚ_SELECT_ERRORrK   rM   rN   rO   ÚEAGAINr&   rP   )ra   r)   ÚrdrF   r8   r:   r:   r;   Úsocket_closed  s(    zSocketChecker.socket_closedN)Ú__name__Ú
__module__Ú__qualname__rb   rl   r:   r:   r:   r;   r[     s   r[   )TNNFNNNFNNFFN).Ú__doc__r    rN   rf   Ústructr]   r\   rd   r   ÚPOLLINÚPOLLPRIÚPOLLERRÚPOLLHUPÚImportErrorr   ri   rL   Zbsonr   Zbson.py3compatr   Zpymongor   r   Zpymongo.commonr   Zpymongo.compression_supportr	   r
   Zpymongo.errorsr   r   r   r   Zpymongo.messager   ÚStructÚunpackr@   r<   rB   r%   rA   rM   Úobjectr[   r:   r:   r:   r;   Ú<module>   sR   

         
%
	