o
    Åîï]Ý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(   ú>/var/www/html/env/lib/python3.10/site-packages/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 |¡ 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'   Úoriginal_indexÚ	operationr(   r(   r)   ÚaddR   s   z_Run.addN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r*   r,   r0   r(   r(   r(   r)   r"   ?   s
    r"   c                 C   sV  |  dd¡}| jtkr|d  |7  < nZ| jtkr"|d  |7  < nL| jtkrn|  d¡}|r\t|ƒ}|D ]}|  |d | ¡|d< q4|d  |¡ |d  |7  < |d  || 7  < n|d  |7  < |d	  |d	 7  < |  d
¡}|r™|D ]!}| ¡ }	|d | }
|  |
¡|	d< | j	|
 |	t
< |d
  |	¡ qw|  d¡}|r©|d  |¡ dS 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ÚresultÚaffectedr8   Ú
n_upsertedÚdocÚwrite_errorsÚreplacementr+   Úwc_errorr(   r(   r)   Ú_merge_command^   s8   





ÿrN   c                 C   s$   | d r| d j dd„ d� t| ƒ‚)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>)Úkey)Úsortr   )rE   r(   r(   r)   Ú_raise_bulk_write_error†   s
   ÿrS   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)Úunicode_decode_error_handlerÚdocument_class)Úcodec_optionsFTN)Úwith_optionsrX   Ú_replaceÚdictÚ
collectionÚorderedr%   ÚexecutedÚbypass_doc_valÚuses_collationÚuses_array_filtersÚis_retryableÚretryingÚstarted_retryable_writeÚcurrent_run©r'   r\   r]   Úbypass_document_validationr(   r(   r)   r*   ’   s    þÿ
z_Bulk.__init__c                 C   s   | j jjj}|r|jstS tS ©N)r\   ÚdatabaseÚclientÚ
_encrypterÚ_bypass_auto_encryptionr   r   )r'   Ú	encrypterr(   r(   r)   Úbulk_ctx_class¥   s   
z_Bulk.bulk_ctx_classc                 C   s:   t d|ƒ t|tƒsd|v stƒ |d< | j t|f¡ dS )z3Add an insert document to the list of ops.
        ÚdocumentÚ_idN)r   Ú
isinstancer   r   r%   r-   r   ©r'   ro   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ur#d| _||d< |dur.d| _||d< |r3d	| _| j t|f¡ dS )
zACreate an update document and add it to the list of ops.
        ÚqÚuÚmultiÚupsertNTÚ	collationÚarrayFiltersF)	r
   r   r   r`   ra   rb   r%   r-   r   )r'   Úselectorr   rv   rw   rx   Ú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ur!d| _||d< | j t|f¡ dS )zACreate a replace document and add it to the list of ops.
        rt   ru   )rv   Frw   NTrx   )r	   r   r   r`   r%   r-   r   )r'   rz   rL   rw   rx   r|   r(   r(   r)   Úadd_replaceÉ   s   ÿz_Bulk.add_replacec                 C   sT   t d|fd|fgƒ}t|ƒ}|durd| _||d< |tkr d| _| j t|f¡ dS )z@Create a delete document and add it to the list of ops.
        rt   ÚlimitNTrx   F)r   r   r`   Ú_DELETE_ALLrb   r%   r-   r   )r'   rz   r   rx   r|   r(   r(   r)   Ú
add_deleteÖ   s   z_Bulk.add_deletec                 c   s^   � d}t | jƒD ]!\}\}}|du rt|ƒ}n|j|kr#|V  t|ƒ}| ||¡ q|V  dS )ziGenerate batches of operations, batched by type of
        operation, in the order **provided**.
        N)Ú	enumerater%   r"   r#   r0   )r'   rD   r+   r#   r/   r(   r(   r)   Úgen_orderedã   s   €


z_Bulk.gen_orderedc                 c   sZ   � t tƒt tƒt tƒg}t| jƒD ]\}\}}||  ||¡ q|D ]}|jr*|V  q"dS )zbGenerate batches of operations, batched by type of
        operation, in arbitrary order.
        N)r"   r   r   r   r‚   r%   r0   )r'   Ú
operationsr+   r#   r/   rD   r(   r(   r)   Úgen_unorderedñ   s   €€þz_Bulk.gen_unorderedc              
   C   sú  |j dk r| jrtdƒ‚|j dk r| jrtdƒ‚| jjj}| jjj}	|	j}
| j	s-t
|ƒ| _	| j	}| |	|¡ |rûtt|j | jjfd| jfgƒ}|jsP|j|d< | jr\|j dkr\d|d	< |  |||||
||j| jj¡}|jt|jƒk ræ|r‰|r�| js�| ¡  d| _| ||tj¡ | |||	¡ t|j|jd ƒ}| ||	¡\}}|  d
i ¡}|  dd¡t!v r¿t" #|¡}t$|||j|ƒ t%|ƒ t$|||j|ƒ d| _&d| _| jrÕd|v rÕn| jt|ƒ7  _|jt|jƒk ss| jrï|d rïd S t
|d ƒ | _	}|s8d S d S )Né   z5Must be connected to MongoDB 3.4+ to use a collation.é   z6Must be connected to MongoDB 3.6+ to use arrayFilters.r]   ÚwriteConcerné   TÚbypassDocumentValidationr=   Úcoder   Fr<   )'Úmax_wire_versionr`   r   ra   r\   ri   Únamerj   Ú_event_listenersre   ÚnextÚvalidate_sessionr   Ú	_COMMANDSr#   r]   Úis_server_defaultro   r_   rn   rX   r&   r@   r%   rd   Ú_start_retryable_writeÚ	_apply_tor   ÚPRIMARYÚsend_cluster_timer   Úexecuter?   r   rB   ÚdeepcopyrN   rS   rc   )r'   Ú	generatorÚwrite_concernÚsessionÚ	sock_infoÚop_idÚ	retryablerE   Údb_namerj   Ú	listenersrD   r|   Úbwcr%   rG   Úto_sendÚwceÚfullr(   r(   r)   Ú_execute_commandý   sh   ÿÿ


ÿ

þ

ã!Ñz_Bulk._execute_commandc              	      s’   g g dddddg dœ‰ t ƒ ‰‡ ‡‡‡‡fdd„}ˆjjj}| |¡�}| ˆj||ˆ¡ W d  ƒ n1 s6w   Y  ˆ d sCˆ d rGtˆ ƒ ˆ S )z&Execute using write commands.
        r   ©r<   r>   r6   r9   r:   r;   r7   r8   c              	      s   ˆ  ˆˆ| |ˆ|ˆ ¡ d S rh   )r¥   )r›   rœ   rž   ©rE   r™   r�   r'   rš   r(   r)   Úretryable_bulkR  s   
þz-_Bulk.execute_command.<locals>.retryable_bulkNr<   r>   )r   r\   ri   rj   Ú_tmp_sessionÚ_retry_with_sessionrb   rS   )r'   r™   rš   r›   r¨   rj   Úsr(   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r$|jdkr$d|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   r]   Úwrˆ   r‰   TrŠ   N)r   r\   r�   r]   Úintr_   rŒ   ri   r   rj   rŽ   r   rX   r   Ú	full_namer%   )	r'   rœ   rD   r�   ÚacknowledgedÚcommandÚ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st|ƒ| _| j}|rqtt	|j
 | j jfddddifgƒ}|  |||||d|j
| j j¡}	|jt|jƒk ret|j|jdƒ}
|	 |
|¡}| jt|ƒ7  _|jt|jƒk sFt|dƒ | _}|sdS dS )zLExecute write commands with OP_MSG and w=0 writeConcern, unordered.
        )r]   Frˆ   r­   r   N)r\   ri   r�   rj   rŽ   r   re   r�   r   r‘   r#   rn   rX   r&   r@   r%   r   Úexecute_unack)r'   rœ   r™   rŸ   rj   r    r�   rD   r|   r¡   r%   r¢   r(   r(   r)   Úexecute_op_msg_no_resultsr  s.   



þ
þüóz_Bulk.execute_op_msg_no_resultsc              	   C   sT   g g dddddg dœ}t ƒ }tƒ }z|  ||d||d|¡ W dS  ty)   Y dS w )zJExecute write commands with OP_MSG and w=0 WriteConcern, ordered.
        r   r¦   NF)r   r   r¥   r   )r'   rœ   r™   rE   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r|jdkrtdƒ‚|jdkr.| jr(|  ||¡S |  ||¡S | j	}t
t| jƒd�}tƒ }t|ƒ}|rÀ|}t|dƒ}| joO|du}z\|jtkr_|  ||||¡ nL|jtkr•|jD ],}	|	d }
d	}|
r|tt|
ƒƒ d
¡r|d}|j||	d |
|	d ||	d ||| j| jd�
 qgn|jD ]}	| ||	d |	d  ||| j¡ q˜W n ty»   | jr¹Y dS Y nw |sBdS dS )z<Execute all operations, returning no results (w=0).
        z3Collation is unsupported for unacknowledged writes.z6arrayFilters is unsupported for unacknowledged writes.r‰   zGCannot set bypass_document_validation with unacknowledged write concernr†   )r­   Nru   Tú$Frt   rw   rv   )rš   r�   r]   r_   r   )r`   r   ra   r_   rŒ   r   r]   r·   r¶   r\   r   r®   r   r�   r#   r   r´   r   r%   ÚiterÚ
startswithÚ_updateÚ_delete)r'   rœ   r™   Úcollrš   r�   Únext_runrD   Ú	needs_ackr/   rJ   Ú
check_keysr(   r(   r)   Úexecute_no_results¦  sz   ÿÿ


ÿ

öû
û€ÿÿÜz_Bulk.execute_no_resultsc                 C   sª   | j stdƒ‚| jrtdƒ‚d| _|p| jj}t||ƒ}| jr$|  ¡ }n|  ¡ }| jj	j
}|jsN| |¡�}|  ||¡ W d  ƒ dS 1 sGw   Y  dS |  |||¡S )zExecute operations.
        zNo operations to executez*Bulk operations can only be executed once.TN)r%   r   r^   r\   rš   r   r]   rƒ   r…   ri   rj   r°   Ú_socket_for_writesrÁ   r¬   )r'   rš   r›   r™   rj   rœ   r(   r(   r)   r—   é  s    


"ÿz_Bulk.execute)FFNN)FNrh   )r1   r2   r3   r4   r*   Úpropertyrn   rs   r}   r~   r�   rƒ   r…   r¥   r¬   r´   r¶   r·   rÁ   r—   r(   r(   r(   r)   rT   �   s,    
	
ÿ
ÿ
ECrT   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   ó   || _ || _|| _d S rh   )Ú_BulkUpsertOperation__selectorÚ_BulkUpsertOperation__bulkÚ_BulkUpsertOperation__collation©r'   rz   Úbulkrx   r(   r(   r)   r*     ó   
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©rv   rw   rx   N©rË   r}   rÊ   rÌ   ©r'   r   r(   r(   r)   Ú
update_one  ó   

þ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
        TrÐ   NrÑ   rÒ   r(   r(   r)   r     rÔ   zBulkUpsertOperation.updatec                 C   ó   | j j| j|d| jd� dS )ú•Replace one entire document matching the selector criteria.

        :Parameters:
          - `replacement` (dict): the replacement document
        T)rw   rx   N)rË   r~   rÊ   rÌ   ©r'   rL   r(   r(   r)   Úreplace_one!  ó   
ÿzBulkUpsertOperation.replace_oneN)	r1   r2   r3   r4   Ú	__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Å   c                 C   rÉ   rh   )Ú_BulkWriteOperation__selectorÚ_BulkWriteOperation__bulkÚ_BulkWriteOperation__collationrÍ   r(   r(   r)   r*   1  rÏ   zBulkWriteOperation.__init__c                 C   rÕ   )zŽUpdate one document matching the selector criteria.

        :Parameters:
          - `update` (dict): the update operations to apply
        F©rv   rx   N©rÝ   r}   rÜ   rÞ   rÒ   r(   r(   r)   rÓ   6  rÙ   zBulkWriteOperation.update_onec                 C   rÕ   )z�Update all documents matching the selector criteria.

        :Parameters:
          - `update` (dict): the update operations to apply
        Trß   Nrà   rÒ   r(   r(   r)   r   ?  rÙ   zBulkWriteOperation.updatec                 C   s   | j j| j|| jd� dS )rÖ   ©rx   N)rÝ   r~   rÜ   rÞ   r×   r(   r(   r)   rØ   H  s   
ÿzBulkWriteOperation.replace_onec                 C   ó   | j j| jt| jd� dS )zARemove a single document matching the selector criteria.
        rá   N)rÝ   r�   rÜ   Ú_DELETE_ONErÞ   ©r'   r(   r(   r)   Ú
remove_oneQ  ó   
ÿzBulkWriteOperation.remove_onec                 C   râ   )z=Remove all documents matching the selector criteria.
        rá   N)rÝ   r�   rÜ   r€   rÞ   rä   r(   r(   r)   ÚremoveW  ræ   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)   rw   ]  s   
ÿzBulkWriteOperation.upsertN)r1   r2   r3   r4   rÚ   r*   rÓ   r   rØ   rå   rç   rw   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)rT   Ú_BulkOperationBuilder__bulkrf   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.

        rz   )r   rÛ   ré   )r'   rz   rx   r(   r(   r)   Úfind‹  s   
zBulkOperationBuilder.findc                 C   s   | j  |¡ dS )zšInsert a single document.

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

        .. seealso:: :ref:`writes-and-ids`
        N)ré   rs   rr   r(   r(   r)   r      s   zBulkOperationBuilder.insertc                 C   s&   |durt di |¤Ž}| 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)   r—   ª  s   zBulkOperationBuilder.execute)TFrh   )	r1   r2   r3   r4   rÚ   r*   rê   r   r—   r(   r(   r(   r)   rè   i  s    
ÿ

rè   )5r4   rB   Ú	itertoolsr   Úbson.objectidr   Úbson.raw_bsonr   Úbson.sonr   Úpymongo.client_sessionr   Úpymongo.commonr   r   r	   r
   Úpymongo.helpersr   Úpymongo.collationr   Úpymongo.errorsr   r   r   r   Úpymongo.messager   r   r   r   r   r   r   Úpymongo.read_preferencesr   Úpymongo.write_concernr   r€   rã   Ú
_BAD_VALUEÚ_UNKNOWN_ERRORÚ_WRITE_CONCERN_ERRORr‘   rC   Úobjectr"   rN   rS   rT   rÄ   rÛ   rè   r(   r(   r(   r)   Ú<module>   s<   $(	  u)>