3
]ð]é*  ã               @   sl   d Z ddlmZ ddlmZ ddlmZmZmZm	Z	 ddl
mZmZmZ G dd„ deƒZG dd	„ d	eƒZd
S )z4CommandCursor class to iterate over command results.é    )Údeque)Úinteger_types)ÚConnectionFailureÚInvalidOperationÚNotMasterErrorÚOperationFailure)Ú_CursorAddressÚ_GetMoreÚ_RawBatchGetMorec               @   sÒ   e Zd ZdZeZd-dd„Zdd„ Zd.d	d
„Zdd„ Z	dd„ Z
dd„ Zdd„ Zedd„ ƒZdd„ Zd/dd„Zdd„ Zedd„ ƒZedd„ ƒZedd „ ƒZed!d"„ ƒZd#d$„ Zd%d&„ ZeZd'd(„ Zd)d*„ Zd+d,„ ZdS )0ÚCommandCursorz)A cursor / iterator over command cursors.r   NFc	       	      C   sª   || _ |d | _t|d ƒ| _|jdƒ| _|| _|| _|| _|| _	|| _
| jdk| _| jrd| jdƒ d|krx|d | _n|j| _| j|ƒ t|tƒ r¦|dk	r¦tdƒ‚dS )	zSCreate a new command cursor.

        The parameter 'retrieved' is unused.
        ÚidÚ
firstBatchÚpostBatchResumeTokenr   TÚnsNz,max_await_time_ms must be an integer or None)Ú_CommandCursor__collectionÚ_CommandCursor__idr   Ú_CommandCursor__dataÚgetÚ$_CommandCursor__postbatchresumetokenÚ_CommandCursor__addressÚ_CommandCursor__batch_sizeÚ!_CommandCursor__max_await_time_msÚ_CommandCursor__sessionÚ _CommandCursor__explicit_sessionÚ_CommandCursor__killedÚ_CommandCursor__end_sessionÚ_CommandCursor__nsÚ	full_nameÚ
batch_sizeÚ
isinstancer   Ú	TypeError)	ÚselfÚ
collectionÚcursor_infoÚaddressÚ	retrievedr   Úmax_await_time_msÚsessionÚexplicit_session© r)   ú9/tmp/pip-build-20mum3z4/pymongo/pymongo/command_cursor.pyÚ__init__!   s&    


zCommandCursor.__init__c             C   s   | j r| j r| jƒ  d S )N)r   r   Ú_CommandCursor__die)r!   r)   r)   r*   Ú__del__@   s    zCommandCursor.__del__c             C   sj   | j }d| _ | jr\| r\t| j| jjƒ}|rH| jjjj| j|| j	d� n| jjjj
| j|ƒ | j|ƒ dS )zCloses this cursor.
        T)r'   N)r   r   r   r   r   r   ÚdatabaseÚclientZ_close_cursor_nowr   Z_close_cursorr   )r!   ÚsynchronousZalready_killedr$   r)   r)   r*   Z__dieD   s    


zCommandCursor.__diec             C   s&   | j r"| j r"| j j|d� d | _ d S )N)Úlock)r   r   Z_end_session)r!   r0   r)   r)   r*   Z__end_sessionU   s    zCommandCursor.__end_sessionc             C   s   | j dƒ dS )z-Explicitly close / kill this cursor.
        TN)r,   )r!   r)   r)   r*   ÚcloseZ   s    zCommandCursor.closec             C   s8   t |tƒstdƒ‚|dk r"tdƒ‚|dkr.dp0|| _| S )aÅ  Limits the number of documents returned in one batch. Each batch
        requires a round trip to the server. It can be adjusted to optimize
        performance and limit data transfer.

        .. note:: batch_size can not override MongoDB's internal limits on the
           amount of data it will return to the client in a single batch (i.e
           if you set batch size to 1,000,000,000, MongoDB will currently only
           return 4-16MB of results per batch).

        Raises :exc:`TypeError` if `batch_size` is not an integer.
        Raises :exc:`ValueError` if `batch_size` is less than ``0``.

        :Parameters:
          - `batch_size`: The size of each batch of results requested.
        zbatch_size must be an integerr   zbatch_size must be >= 0é   é   )r   r   r    Ú
ValueErrorr   )r!   r   r)   r)   r*   r   _   s    
zCommandCursor.batch_sizec             C   s   t | jƒdkS )zUReturns `True` if the cursor has documents remaining from the
        previous batch.r   )Úlenr   )r!   r)   r)   r*   Ú	_has_nextw   s    zCommandCursor._has_nextc             C   s   | j S )zcRetrieve the postBatchResumeToken from the response to a
        changeStream aggregate or getMore.)r   )r!   r)   r)   r*   Ú_post_batch_resume_token|   s    z&CommandCursor._post_batch_resume_tokenc       
         s  ‡ fdd„}ˆ j jj}y|j|ˆ jˆ jd�}W nl tk
rJ   |ƒ  ‚ Y nR tk
rd   |ƒ  ‚ Y n8 tk
r~   |ƒ  ‚ Y n t	k
rš   ˆ j
ƒ  ‚ Y nX |j}|j}|j}|rÞ|d d }|d }	|jdƒˆ _|d ˆ _n|}	|jˆ _ˆ jdkrú|ƒ  t|	ƒˆ _d	S )
z8Send a getmore message and handle the response.
        c                  s   dˆ _ ˆ jdƒ d S )NT)r   r   r)   )r!   r)   r*   Úkill…   s    z*CommandCursor.__send_message.<locals>.kill)r$   r   ÚcursorZ	nextBatchr   r   N)r   r.   r/   Z_run_operation_with_responseÚ_unpack_responser   r   r   r   Ú	Exceptionr,   Úfrom_commandÚdataÚdocsr   r   r   Ú	cursor_idr   r   )
r!   Z	operationr9   r/   Úresponser=   Zreplyr?   r:   Z	documentsr)   )r!   r*   Z__send_message‚   s<    

zCommandCursor.__send_messagec             C   s   |j ||||ƒS )N)Zunpack_response)r!   rA   r@   Úcodec_optionsÚuser_fieldsÚlegacy_responser)   r)   r*   r;   ²   s    
zCommandCursor._unpack_responsec             C   s�   t | jƒs| jrt | jƒS | jrv| jjddƒ\}}| jj| jƒ}| j	| j
||| j| j| jj|| j| jjj| jdƒ
ƒ nd| _| jdƒ t | jƒS )a  Refreshes the cursor with more data from the server.

        Returns the length of self.__data after refresh. Will exit early if
        self.__data is already non-empty. Raises OperationFailure when the
        cursor cannot be refreshed due to an error on the query.
        Ú.r3   FT)r6   r   r   r   r   Úsplitr   Z_read_preference_forr'   Ú_CommandCursor__send_messageÚ_getmore_classr   rB   r   r.   r/   r   r   )r!   ZdbnameZcollnameZ	read_prefr)   r)   r*   Ú_refresh·   s&    


zCommandCursor._refreshc             C   s   t t| jƒp| j ƒS )a  Does this cursor have the potential to return more data?

        Even if :attr:`alive` is ``True``, :meth:`next` can raise
        :exc:`StopIteration`. Best to use a for loop::

            for doc in collection.aggregate(pipeline):
                print(doc)

        .. note:: :attr:`alive` can be True while iterating a cursor from
          a failed server. In this case :attr:`alive` will return False after
          :meth:`next` fails to retrieve the next batch of results from the
          server.
        )Úboolr6   r   r   )r!   r)   r)   r*   ÚaliveÕ   s    zCommandCursor.alivec             C   s   | j S )zReturns the id of the cursor.)r   )r!   r)   r)   r*   r@   æ   s    zCommandCursor.cursor_idc             C   s   | j S )zUThe (host, port) of the server used, or None.

        .. versionadded:: 3.0
        )r   )r!   r)   r)   r*   r$   ë   s    zCommandCursor.addressc             C   s   | j r| jS dS )zmThe cursor's :class:`~pymongo.client_session.ClientSession`, or None.

        .. versionadded:: 3.6
        N)r   r   )r!   r)   r)   r*   r'   ó   s    zCommandCursor.sessionc             C   s   | S )Nr)   )r!   r)   r)   r*   Ú__iter__ü   s    zCommandCursor.__iter__c             C   s*   x | j r | jdƒ}|dk	r|S qW t‚dS )zAdvance the cursor.TN)rK   Ú	_try_nextÚStopIteration)r!   Údocr)   r)   r*   Únextÿ   s
    
zCommandCursor.nextc             C   sL   t | jƒ r | j r |r | jƒ  t | jƒrD| j}|jj| jjƒ |ƒS dS dS )z<Advance the cursor blocking for at most one getMore command.N)r6   r   r   rI   r   r.   Z_fix_outgoingÚpopleft)r!   Zget_more_allowedZcollr)   r)   r*   rM     s    
zCommandCursor._try_nextc             C   s   | S )Nr)   )r!   r)   r)   r*   Ú	__enter__  s    zCommandCursor.__enter__c             C   s   | j ƒ  d S )N)r2   )r!   Úexc_typeÚexc_valÚexc_tbr)   r)   r*   Ú__exit__  s    zCommandCursor.__exit__)r   r   NNF)F)NF)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r	   rH   r+   r-   r,   r   r2   r   r7   Úpropertyr8   rG   r;   rI   rK   r@   r$   r'   rL   rP   Ú__next__rM   rR   rV   r)   r)   r)   r*   r      s2     

1
	

r   c                   s4   e Zd ZeZd
‡ fdd„	Zddd„Zdd	„ Z‡  ZS )ÚRawBatchCommandCursorr   NFc	       	   	      s2   |j dƒ st‚tt| ƒj||||||||ƒ dS )a  Create a new cursor / iterator over raw batches of BSON data.

        Should not be called directly by application developers -
        see :meth:`~pymongo.collection.Collection.aggregate_raw_batches`
        instead.

        .. mongodoc:: cursors
        r   N)r   ÚAssertionErrorÚsuperr]   r+   )	r!   r"   r#   r$   r%   r   r&   r'   r(   )Ú	__class__r)   r*   r+     s    

zRawBatchCommandCursor.__init__c             C   s
   |j |ƒS )N)Zraw_response)r!   rA   r@   rB   rC   rD   r)   r)   r*   r;   /  s    z&RawBatchCommandCursor._unpack_responsec             C   s   t dƒ‚d S )Nz)Cannot call __getitem__ on RawBatchCursor)r   )r!   Úindexr)   r)   r*   Ú__getitem__3  s    z!RawBatchCommandCursor.__getitem__)r   r   NNF)NF)	rW   rX   rY   r
   rH   r+   r;   rb   Ú__classcell__r)   r)   )r`   r*   r]     s     
r]   N)rZ   Úcollectionsr   Zbson.py3compatr   Zpymongo.errorsr   r   r   r   Zpymongo.messager   r	   r
   Úobjectr   r]   r)   r)   r)   r*   Ú<module>   s     