3
]ð]Ýf  ã               @   sL  d Z ddlZddlmZ ddlmZ ddlm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 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mZmZm Z m!Z! ddl"m#Z# ddl$m%Z% dZ&dZ'dZ(dZ)dZ*d%Z+dZ,G dd„ de-ƒZ.dd„ Z/dd„ Z0G dd„ de-ƒZ1G dd „ d e-ƒZ2G d!d"„ d"e-ƒZ3G d#d$„ d$e-ƒZ4dS )&z<The bulk write operations interface.

.. versionadded:: 2.7
é    N)Úislice)ÚObjectId)ÚRawBSONDocument)ÚSON)Ú_validate_session_write_concern)Úvalidate_is_mappingÚvalidate_is_document_typeÚvalidate_ok_for_replaceÚvalidate_ok_for_update)Ú_RETRYABLE_ERROR_CODES)Úvalidate_collation_or_none)ÚBulkWriteErrorÚConfigurationErrorÚInvalidOperationÚOperationFailure)Ú_INSERTÚ_UPDATEÚ_DELETEÚ_do_batched_insertÚ_randintÚ_BulkWriteContextÚ_EncryptedBulkWriteContext)ÚReadPreference)ÚWriteConcerné   é   é   é@   ÚinsertÚupdateÚdeleteÚopc               @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	Ú_Runz,Represents a batch of write operations.
    c             C   s   || _ g | _g | _d| _dS )z%Initialize a new Run object.
        r   N)Úop_typeÚ	index_mapÚopsÚ
idx_offset)Úselfr#   © r(   ú//tmp/pip-build-20mum3z4/pymongo/pymongo/bulk.pyÚ__init__B   s    z_Run.__init__c             C   s
   | j | S )z”Get the original index of an operation in this run.

        :Parameters:
          - `idx`: The Run index that maps to the original index.
        )r$   )r'   Úidxr(   r(   r)   ÚindexJ   s    z
_Run.indexc             C   s   | j j|ƒ | jj|ƒ dS )zåAdd an operation to this Run instance.

        :Parameters:
          - `original_index`: The original index of this operation
            within a larger bulk operation.
          - `operation`: The operation document.
        N)r$   Úappendr%   )r'   Zoriginal_indexÚ	operationr(   r(   r)   ÚaddR   s    z_Run.addN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r*   r,   r/   r(   r(   r(   r)   r"   ?   s   r"   c             C   s^  |j ddƒ}| jtkr(|d  |7  < n¸| jtkrD|d  |7  < nœ| jtkrà|j dƒ}|r¼t|ƒ}x"|D ]}| j|d | ƒ|d< qjW |d j|ƒ |d  |7  < |d  || 7  < n|d  |7  < |d	  |d	 7  < |j d
ƒ}|�r<xJ|D ]B}|jƒ }	|d | }
| j|
ƒ|	d< | j	|
 |	t
< |d
 j|	ƒ qöW |j dƒ}|�rZ|d j|ƒ dS )z<Merge a write command result into the full bulk result.
    Únr   Ú	nInsertedÚnRemovedÚupsertedr,   Ú	nUpsertedÚnMatchedÚ	nModifiedÚwriteErrorsÚwriteConcernErrorÚwriteConcernErrorsN)Úgetr#   r   r   r   Úlenr,   ÚextendÚcopyr%   Ú_UOPr-   )ÚrunÚfull_resultÚoffsetÚresultZaffectedr7   Z
n_upsertedÚdocZwrite_errorsÚreplacementr+   Zwc_errorr(   r(   r)   Ú_merge_command^   s6    







rI   c             C   s(   | d r| d j dd„ d� t| ƒ‚dS )z:Raise a BulkWriteError from the full bulk api result.
    r;   c             S   s   | d S )Nr,   r(   )Úerrorr(   r(   r)   Ú<lambda>‹   s    z)_raise_bulk_write_error.<locals>.<lambda>)ÚkeyN)Úsortr   )rD   r(   r(   r)   Ú_raise_bulk_write_error†   s    rN   c               @   s’   e Zd ZdZdd„ Zedd„ ƒZdd„ Zd"d
d„Zd#dd„Z	d$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	S )%Ú_Bulkz,The private guts of the bulk write API.
    c             C   sZ   |j |jjdtd�d�| _|| _g | _d| _|| _d| _	d| _
d| _d| _d| _d| _dS )z%Initialize a _Bulk instance.
        Úreplace)Zunicode_decode_error_handlerZdocument_class)Úcodec_optionsFTN)Zwith_optionsrQ   Ú_replaceÚdictÚ
collectionÚorderedr%   ÚexecutedÚbypass_doc_valÚuses_collationÚuses_array_filtersÚis_retryableÚretryingÚstarted_retryable_writeÚcurrent_run)r'   rT   rU   Úbypass_document_validationr(   r(   r)   r*   ’   s    z_Bulk.__init__c             C   s$   | j jjj}|r|j rtS tS d S )N)rT   ÚdatabaseÚclientZ
_encrypterZ_bypass_auto_encryptionr   r   )r'   Z	encrypterr(   r(   r)   Úbulk_ctx_class¥   s    z_Bulk.bulk_ctx_classc             C   s:   t d|ƒ t|tƒpd|ks&tƒ |d< | jjt|fƒ dS )z3Add an insert document to the list of ops.
        ÚdocumentÚ_idN)r   Ú
isinstancer   r   r%   r-   r   )r'   rb   r(   r(   r)   Ú
add_insert­   s    

z_Bulk.add_insertFNc             C   sz   t |ƒ td|fd|fd|fd|fgƒ}t|ƒ}|dk	rFd| _||d< |dk	r\d| _||d< |rfd	| _| jjt|fƒ dS )
zACreate an update document and add it to the list of ops.
        ÚqÚuÚmultiÚupsertNTÚ	collationZarrayFiltersF)	r
   r   r   rX   rY   rZ   r%   r-   r   )r'   Úselectorr   rh   ri   rj   Zarray_filtersÚcmdr(   r(   r)   Ú
add_update¶   s    z_Bulk.add_updatec             C   sV   t |ƒ td|fd|fd	d|fgƒ}t|ƒ}|dk	rBd| _||d< | jjt|fƒ dS )
zACreate a replace document and add it to the list of ops.
        rf   rg   rh   Fri   NTrj   )rh   F)r	   r   r   rX   r%   r-   r   )r'   rk   rH   ri   rj   rl   r(   r(   r)   Úadd_replaceÉ   s    z_Bulk.add_replacec             C   sT   t d|fd|fgƒ}t|ƒ}|dk	r2d| _||d< |tkr@d| _| jjt|fƒ dS )z@Create a delete document and add it to the list of ops.
        rf   ÚlimitNTrj   F)r   r   rX   Ú_DELETE_ALLrZ   r%   r-   r   )r'   rk   ro   rj   rl   r(   r(   r)   Ú
add_deleteÖ   s    z_Bulk.add_deletec             c   s`   d}xPt | jƒD ]B\}\}}|dkr.t|ƒ}n|j|krF|V  t|ƒ}|j||ƒ qW |V  dS )ziGenerate batches of operations, batched by type of
        operation, in the order **provided**.
        N)Ú	enumerater%   r"   r#   r/   )r'   rC   r+   r#   r.   r(   r(   r)   Úgen_orderedã   s    

z_Bulk.gen_orderedc             c   s`   t tƒt tƒt tƒg}x*t| jƒD ]\}\}}|| j||ƒ q"W x|D ]}|jrH|V  qHW dS )zbGenerate batches of operations, batched by type of
        operation, in arbitrary order.
        N)r"   r   r   r   rr   r%   r/   )r'   Ú
operationsr+   r#   r.   rC   r(   r(   r)   Úgen_unorderedñ   s    
z_Bulk.gen_unorderedc          	   C   s  |j dk r| jrtdƒ‚|j dk r0| jr0tdƒ‚| jjj}| jjj}	|	j}
| j	sZt
|ƒ| _	| j	}|j|	|ƒ �x�|�rþtt|j | jjfd| jfgƒ}|js¦|j|d< | jr¾|j dkr¾d|d	< | j|||||
||j| jjƒ}xú|jt|jƒk �rÖ|�r$|�r| j �r|jƒ  d| _|j||tjƒ |j|||	ƒ t|j|jd ƒ}|j||	ƒ\}}|j d
i ƒ}|j ddƒt!k�r’t"j#|ƒ}t$|||j|ƒ t%|ƒ t$|||j|ƒ d| _&d| _| j�rÂd|k�rÂP | jt|ƒ7  _qÞW | j�rì|d �rìP t
|d ƒ | _	}qpW d S )Né   z5Must be connected to MongoDB 3.4+ to use a collation.é   z6Must be connected to MongoDB 3.6+ to use arrayFilters.rU   ÚwriteConcerné   TÚbypassDocumentValidationr<   Úcoder   Fr;   )'Úmax_wire_versionrX   r   rY   rT   r_   Únamer`   Ú_event_listenersr]   ÚnextZvalidate_sessionr   Ú	_COMMANDSr#   rU   Zis_server_defaultrb   rW   ra   rQ   r&   r?   r%   r\   Z_start_retryable_writeZ	_apply_tor   ZPRIMARYZsend_cluster_timer   Úexecuter>   r   rA   ÚdeepcopyrI   rN   r[   )r'   Ú	generatorÚwrite_concernÚsessionÚ	sock_infoÚop_idÚ	retryablerD   Údb_namer`   Ú	listenersrC   rl   Úbwcr%   rF   Úto_sendZwceÚfullr(   r(   r)   Ú_execute_commandý   s\    





z_Bulk._execute_commandc                s~   g g dddddg dœ‰ t ƒ ‰‡ ‡‡‡‡fdd„}ˆjjj}|j|ƒ�}|jˆj||ˆƒ W dQ R X ˆ d srˆ d rztˆ ƒ ˆ S )z&Execute using write commands.
        r   )r;   r=   r5   r8   r9   r:   r6   r7   c                s   ˆj ˆˆ| |ˆ|ˆ ƒ d S )N)rŽ   )r…   r†   rˆ   )rD   rƒ   r‡   r'   r„   r(   r)   Úretryable_bulkR  s    
z-_Bulk.execute_command.<locals>.retryable_bulkNr;   r=   )r   rT   r_   r`   Z_tmp_sessionZ_retry_with_sessionrZ   rN   )r'   rƒ   r„   r…   r�   r`   Úsr(   )rD   rƒ   r‡   r'   r„   r)   Úexecute_commandB  s"    
z_Bulk.execute_commandc       	   	   C   s˜   t d| jjfd| jfgƒ}dt| jƒi}||d< | jrH|jdkrHd|d< | jj}t|j||||j	j
dt| jjƒ}t| jj|jd||| j | jj|ƒ dS )	z.Execute insert, returning no results.
        r   rU   Úwrx   ry   Trz   N)r   rT   r}   rU   ÚintrW   r|   r_   r   r`   r~   r   rQ   r   Z	full_namer%   )	r'   r†   rC   r‡   ÚacknowledgedÚcommandZconcernÚdbr‹   r(   r(   r)   Úexecute_insert_no_results`  s    z_Bulk.execute_insert_no_resultsc          	   C   sÚ   | j jj}| j jj}|j}tƒ }| js0t|ƒ| _| j}xž|rÔtt	|j
 | j jfddddifgƒ}| j|||||d|j
| j jƒ}	xB|jt|jƒk rÀt|j|jdƒ}
|	j|
|ƒ}| jt|ƒ7  _q€W t|dƒ | _}q8W dS )zLExecute write commands with OP_MSG and w=0 writeConcern, unordered.
        rU   Frx   r’   r   N)rU   F)rT   r_   r}   r`   r~   r   r]   r   r   r€   r#   ra   rQ   r&   r?   r%   r   Zexecute_unack)r'   r†   rƒ   r‰   r`   rŠ   r‡   rC   rl   r‹   r%   rŒ   r(   r(   r)   Úexecute_op_msg_no_resultsr  s&    


z_Bulk.execute_op_msg_no_resultsc             C   sV   g g dddddg dœ}t ƒ }tƒ }y| j||d||d|ƒ W n tk
rP   Y nX dS )zJExecute write commands with OP_MSG and w=0 WriteConcern, ordered.
        r   )r;   r=   r5   r8   r9   r:   r6   r7   NF)r   r   rŽ   r   )r'   r†   rƒ   rD   r„   r‡   r(   r(   r)   Úexecute_command_no_results�  s     z _Bulk.execute_command_no_resultsc             C   s–  | j rtdƒ‚| jrtdƒ‚| jr4|jdkr4tdƒ‚|jdkr\| jrP| j||ƒS | j||ƒS | j	}t
t| jƒd�}tƒ }t|ƒ}�x|�r�|}t|dƒ}| jo¤|dk	}yÆ|jtkrÄ| j||||ƒ n¦|jtk�r8x˜|jD ]Z}	|	d }
d	}|
oütt|
ƒƒjd
ƒ�rd}|j||	d |
|	d ||	d ||| j| jd�
 qØW n2x0|jD ]&}	|j||	d |	d  ||| jƒ �q@W W q„ tk
�rŒ   | j�rˆP Y q„X q„W dS )z<Execute all operations, returning no results (w=0).
        z3Collation is unsupported for unacknowledged writes.z6arrayFilters is unsupported for unacknowledged writes.ry   zGCannot set bypass_document_validation with unacknowledged write concernrv   )r’   Nrg   Tú$Frf   ri   rh   )r„   r‡   rU   rW   ro   )rX   r   rY   rW   r|   r   rU   r™   r˜   rT   r   r“   r   r   r#   r   r—   r   r%   ÚiterÚ
startswithÚ_updateÚ_delete)r'   r†   rƒ   Zcollr„   r‡   Znext_runrC   Z	needs_ackr.   rG   Z
check_keysr(   r(   r)   Úexecute_no_results¦  sf    



z_Bulk.execute_no_resultsc          
   C   s–   | j stdƒ‚| jrtdƒ‚d| _|p,| jj}t||ƒ}| jrH| jƒ }n| jƒ }| jj	j
}|js„|j|ƒ�}| j||ƒ W dQ R X n| j|||ƒS dS )zExecute operations.
        zNo operations to executez*Bulk operations can only be executed once.TN)r%   r   rV   rT   r„   r   rU   rs   ru   r_   r`   r”   Z_socket_for_writesrŸ   r‘   )r'   r„   r…   rƒ   r`   r†   r(   r(   r)   r�   é  s    


z_Bulk.execute)FFNN)FN)N)r0   r1   r2   r3   r*   Úpropertyra   re   rm   rn   rq   rs   ru   rŽ   r‘   r—   r˜   r™   rŸ   r�   r(   r(   r(   r)   rO   �   s$   	 
 

ECrO   c               @   s4   e Zd ZdZdZdd„ Zdd„ Zd	d
„ Zdd„ ZdS )ÚBulkUpsertOperationz/An interface for adding upsert operations.
    Ú
__selectorÚ__bulkÚ__collationc             C   s   || _ || _|| _d S )N)Ú_BulkUpsertOperation__selectorÚ_BulkUpsertOperation__bulkÚ_BulkUpsertOperation__collation)r'   rk   Úbulkrj   r(   r(   r)   r*     s    zBulkUpsertOperation.__init__c             C   s   | j j| j|dd| jd� dS )z…Update one document matching the selector.

        :Parameters:
          - `update` (dict): the update operations to apply
        FT)rh   ri   rj   N)r¦   rm   r¥   r§   )r'   r   r(   r(   r)   Ú
update_one  s    
zBulkUpsertOperation.update_onec             C   s   | j j| j|dd| jd� dS )z†Update all documents matching the selector.

        :Parameters:
          - `update` (dict): the update operations to apply
        T)rh   ri   rj   N)r¦   rm   r¥   r§   )r'   r   r(   r(   r)   r     s    
zBulkUpsertOperation.updatec             C   s   | j j| j|d| jd� dS )z•Replace one entire document matching the selector criteria.

        :Parameters:
          - `replacement` (dict): the replacement document
        T)ri   rj   N)r¦   rn   r¥   r§   )r'   rH   r(   r(   r)   Úreplace_one!  s    zBulkUpsertOperation.replace_oneN)r¢   r£   r¤   )	r0   r1   r2   r3   Ú	__slots__r*   r©   r   rª   r(   r(   r(   r)   r¡     s   

r¡   c               @   sL   e Zd 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dS )ÚBulkWriteOperationz9An interface for adding update or remove operations.
    r¢   r£   r¤   c             C   s   || _ || _|| _d S )N)Ú_BulkWriteOperation__selectorÚ_BulkWriteOperation__bulkÚ_BulkWriteOperation__collation)r'   rk   r¨   rj   r(   r(   r)   r*   1  s    zBulkWriteOperation.__init__c             C   s   | j j| j|d| jd� dS )zŽUpdate one document matching the selector criteria.

        :Parameters:
          - `update` (dict): the update operations to apply
        F)rh   rj   N)r®   rm   r­   r¯   )r'   r   r(   r(   r)   r©   6  s    zBulkWriteOperation.update_onec             C   s   | j j| j|d| jd� dS )z�Update all documents matching the selector criteria.

        :Parameters:
          - `update` (dict): the update operations to apply
        T)rh   rj   N)r®   rm   r­   r¯   )r'   r   r(   r(   r)   r   ?  s    zBulkWriteOperation.updatec             C   s   | j j| j|| jd� dS )z•Replace one entire document matching the selector criteria.

        :Parameters:
          - `replacement` (dict): the replacement document
        )rj   N)r®   rn   r­   r¯   )r'   rH   r(   r(   r)   rª   H  s    zBulkWriteOperation.replace_onec             C   s   | j j| jt| jd� dS )zARemove a single document matching the selector criteria.
        )rj   N)r®   rq   r­   Ú_DELETE_ONEr¯   )r'   r(   r(   r)   Ú
remove_oneQ  s    zBulkWriteOperation.remove_onec             C   s   | j j| jt| jd� dS )z=Remove all documents matching the selector criteria.
        )rj   N)r®   rq   r­   rp   r¯   )r'   r(   r(   r)   ÚremoveW  s    zBulkWriteOperation.removec             C   s   t | j| j| jƒS )zØSpecify that all chained update operations should be
        upserts.

        :Returns:
          - A :class:`BulkUpsertOperation` instance, used to add
            update operations to this bulk operation.
        )r¡   r­   r®   r¯   )r'   r(   r(   r)   ri   ]  s    
zBulkWriteOperation.upsertN)r¢   r£   r¤   )r0   r1   r2   r3   r«   r*   r©   r   rª   r±   r²   ri   r(   r(   r(   r)   r¬   +  s   			r¬   c               @   s:   e Zd ZdZdZddd„Zddd	„Zd
d„ Zddd„ZdS )ÚBulkOperationBuilderzL**DEPRECATED**: An interface for executing a batch of write operations.
    r£   TFc             C   s   t |||ƒ| _dS )a(  **DEPRECATED**: Initialize a new BulkOperationBuilder instance.

        :Parameters:
          - `collection`: A :class:`~pymongo.collection.Collection` instance.
          - `ordered` (optional): If ``True`` all operations will be executed
            serially, in the order provided, and the entire execution will
            abort on the first error. If ``False`` operations will be executed
            in arbitrary order (possibly in parallel on the server), reporting
            any errors that occurred after attempting all operations. Defaults
            to ``True``.
          - `bypass_document_validation`: (optional) If ``True``, allows the
            write to opt-out of document level validation. Default is
            ``False``.

        .. note:: `bypass_document_validation` requires server version
          **>= 3.2**

        .. versionchanged:: 3.5
           Deprecated. Use :meth:`~pymongo.collection.Collection.bulk_write`
           instead.

        .. versionchanged:: 3.2
          Added bypass_document_validation support
        N)rO   Ú_BulkOperationBuilder__bulk)r'   rT   rU   r^   r(   r(   r)   r*   o  s    zBulkOperationBuilder.__init__Nc             C   s   t d|ƒ t|| j|ƒS )a;  Specify selection criteria for bulk operations.

        :Parameters:
          - `selector` (dict): the selection criteria for update
            and remove operations.
          - `collation` (optional): An instance of
            :class:`~pymongo.collation.Collation`. This option is only
            supported on MongoDB 3.4 and above.

        :Returns:
          - A :class:`BulkWriteOperation` instance, used to add
            update and remove operations to this bulk operation.

        .. versionchanged:: 3.4
           Added the `collation` option.

        rk   )r   r¬   r´   )r'   rk   rj   r(   r(   r)   Úfind‹  s    
zBulkOperationBuilder.findc             C   s   | j j|ƒ dS )zšInsert a single document.

        :Parameters:
          - `document` (dict): the document to insert

        .. seealso:: :ref:`writes-and-ids`
        N)r´   re   )r'   rb   r(   r(   r)   r      s    zBulkOperationBuilder.insertc             C   s"   |dk	rt f |Ž}| jj|dd�S )zœExecute all provided operations.

        :Parameters:
          - write_concern (optional): the write concern for this bulk
            execution.
        N)r…   )r   r´   r�   )r'   r„   r(   r(   r)   r�   ª  s    
zBulkOperationBuilder.execute)TF)N)N)	r0   r1   r2   r3   r«   r*   rµ   r   r�   r(   r(   r(   r)   r³   i  s    


r³   )r   r   r    )5r3   rA   Ú	itertoolsr   Zbson.objectidr   Zbson.raw_bsonr   Zbson.sonr   Zpymongo.client_sessionr   Zpymongo.commonr   r   r	   r
   Zpymongo.helpersr   Zpymongo.collationr   Zpymongo.errorsr   r   r   r   Zpymongo.messager   r   r   r   r   r   r   Zpymongo.read_preferencesr   Zpymongo.write_concernr   rp   r°   Z
_BAD_VALUEZ_UNKNOWN_ERRORZ_WRITE_CONCERN_ERRORr€   rB   Úobjectr"   rI   rN   rO   r¡   r¬   r³   r(   r(   r(   r)   Ú<module>   s:   $(	  u)>