3
]ð]¸  ã               @   s‚   d Z 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
 ddlmZmZ ddlmZ d	d
d
dœiZG dd„ deƒZdS )z2Communicate with one MongoDB server in a topology.é    )Údatetime)Ú_decode_all_selective)ÚNotMasterErrorÚOperationFailure)Ú_check_command_response)Ú_convert_exception)ÚResponseÚExhaustResponse)ÚSERVER_TYPEÚcursoré   )Ú
firstBatchÚ	nextBatchc               @   s~   e Zd Zddd„Zdd„ Zdd„ Zdd	„ Zd
d„ Zdd„ Zddd„Z	e
dd„ ƒZejdd„ ƒZe
dd„ ƒZdd„ Zdd„ ZdS )ÚServerNc             C   sF   || _ || _|| _|| _|dk	o$|j| _|| _d| _| jrB|ƒ | _dS )zRepresent one MongoDB server.N)Ú_descriptionÚ_poolÚ_monitorÚ_topology_idZenabled_for_serverÚ_publishÚ	_listenerÚ_events)ÚselfÚserver_descriptionÚpoolZmonitorZtopology_idÚ	listenersÚevents© r   ú1/tmp/pip-build-20mum3z4/pymongo/pymongo/server.pyÚ__init__   s    zServer.__init__c             C   s   | j jƒ  dS )z[Start monitoring, or restart after a fork.

        Multiple calls have no effect.
        N)r   Úopen)r   r   r   r   r   ,   s    zServer.openc             C   s   | j jƒ  dS )zClear the connection pool.N)r   Úreset)r   r   r   r   r    3   s    zServer.resetc             C   s<   | j r$| jj| jj| jj| jffƒ | jj	ƒ  | j
jƒ  dS )zXClear the connection pool and stop the monitor.

        Reconnect with open().
        N)r   r   Úputr   Zpublish_server_closedr   Úaddressr   r   Úcloser   r    )r   r   r   r   r#   7   s
    
zServer.closec             C   s   | j jƒ  dS )zCheck the server's state soon.N)r   Úrequest_check)r   r   r   r   r$   B   s    zServer.request_checkc             C   sz  d}|j }|rtjƒ }	|j }
|
rN|j||ƒ}|j|||ƒ}| j|ƒ\}}}nd}d}|r‚|j|ƒ\}}|j||||j	ƒ tjƒ }	yz|
r |j
||ƒ |j|ƒ}n
|jdƒ}|r¸t}d}nd}d}|||j|j||d�}|rú|d }|jj||jƒ t|ƒ W nn tk
�rj } zP|�rXtjƒ |	 }t|ttfƒ�r:|j}nt|ƒ}|j|||j||j	ƒ ‚ W Y dd}~X nX |�r tjƒ |	 }|�rŽ|d }n\|jdk�r®|�r¨|d ni }n<|j|jƒ dœdd	œ}|jd
k�rÞ||d d< n||d d< |j|||j||j	ƒ |j}|�r8|j�r8|�r8|jj|jƒ ƒ}t ||j|ƒ}|�r^t!|| j"j	|| j#||||d�}nt$|| j"j	||||d�}|S )aˆ  Run a _Query or _GetMore operation and return a Response object.

        This method is used only to run _Query/_GetMore operations from
        cursors.
        Can raise ConnectionFailure, OperationFailure, etc.

        :Parameters:
          - `operation`: A _Query or _GetMore object.
          - `set_slave_okay`: Pass to operation.get_message.
          - `all_credentials`: dict, maps auth source to MongoCredential.
          - `listeners`: Instance of _EventListeners or None.
          - `exhaust`: If True, then this is an exhaust cursor operation.
          - `unpack_res`: A callable that decodes the wire protocol response.
        NFr   T)Úlegacy_responseÚuser_fieldsZexplain)ÚidÚnsr   )r   ÚokÚfindr   r   r   )Údatar"   Zsocket_infor   ÚdurationÚ
request_idÚfrom_commandÚdocs)r+   r"   r,   r-   r.   r/   )%Zenabled_for_commandsr   ÚnowZexhaust_mgrZuse_commandZget_messageÚ_split_messageZ
as_commandZpublish_command_startr"   Úsend_messageZreceive_messageÚ_CURSOR_DOC_FIELDSZ	cursor_idZcodec_optionsÚclientZ_process_responseÚsessionr   Ú	ExceptionÚ
isinstancer   r   Údetailsr   Zpublish_command_failureÚnameÚ	namespaceZpublish_command_successZ
_encrypterZdecryptZraw_command_responser   r	   r   r   r   )r   Z	sock_infoZ	operationZset_slave_okayr   ZexhaustZ
unpack_resr,   ÚpublishÚstartr2   Zuse_cmdÚmessager-   r+   Zmax_doc_sizeÚcmdZdbnZreplyr&   r%   r/   ÚfirstÚexcZfailureÚresr4   Z	decryptedÚresponser   r   r   Úrun_operation_with_responseF   s¬    








z"Server.run_operation_with_responseFc             C   s   | j j||ƒS )N)r   Ú
get_socket)r   Zall_credentialsÚcheckoutr   r   r   rD   Ç   s    zServer.get_socketc             C   s   | j S )N)r   )r   r   r   r   ÚdescriptionÊ   s    zServer.descriptionc             C   s   |j | jj kst‚|| _d S )N)r"   r   ÚAssertionError)r   r   r   r   r   rF   Î   s    c             C   s   | j S )N)r   )r   r   r   r   r   Ó   s    zServer.poolc             C   s&   t |ƒdkr|S |\}}||dfS dS )z“Return request_id, data, max_doc_size.

        :Parameters:
          - `message`: (request_id, data, max_doc_size) or (request_id, data)
        é   r   N)Úlen)r   r=   r-   r+   r   r   r   r1   ×   s    zServer._split_messagec             C   s(   | j }d|jd |jd tj|j f S )Nz<Server "%s:%s" %s>r   r   )r   r"   r
   Ú_fieldsZserver_type)r   Údr   r   r   Ú__str__ä   s    zServer.__str__)NNN)F)Ú__name__Ú
__module__Ú__qualname__r   r   r    r#   r$   rC   rD   ÚpropertyrF   Úsetterr   r1   rL   r   r   r   r   r      s    
 
r   N)Ú__doc__r   Zbsonr   Zpymongo.errorsr   r   Zpymongo.helpersr   Zpymongo.messager   Zpymongo.responser   r	   Zpymongo.server_typer
   r3   Úobjectr   r   r   r   r   Ú<module>   s   