3
]ð]é  ã               @   s–  d Z ddlZddlZddlZddlZddlmZmZmZmZm	Z	 ddl
mZ ddlmZmZ ddlmZmZ ddlmZ yddlmZ d	ZW n ek
r¤   d
ZY nX ddlmZ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dl%m&Z& dZ'dÂZ(dZ)dZ*dZ+dZ,dZ-dZ.dZ/dZ0dZ1dZ2dZ3e*de+de,diZ4ddd d!œZ5d"Z6ed#d$�Z7d%d&„ Z8d'd(„ Z9d)d*„ Z:d+d,„ Z;edÃdÄdÅdÆdÇgƒZ<edÈdÉdÊdËdÌdÍdÎdÏdÐdÑdÒdÓgƒZ=dÔdMdN„Z>dOdP„ Z?G dQdR„ dRe@ƒZAG dSdT„ dTe@ƒZBG dUdV„ dVeAƒZCG dWdX„ dXeBƒZDG dYdZ„ dZeEƒZFejGd[ƒjHZId\ZJd]d^„ ZKd_d`„ ZLejGdaƒjHZMdbdc„ ZNejGddƒjHZOdedf„ ZPdgdh„ ZQdidj„ ZRe�rhejSZRdÕdkdl„ZTdmdn„ ZUdodp„ ZVdqdr„ ZWe�r–ejXZWdÖdsdt„ZYejGduƒjHZZejGdvƒjHZ[dwdx„ Z\dydz„ Z]d{d|„ Z^e�rÜej_Z^d×d}d~„Z_dd€„ Z`dØd�d‚„ZadÙdƒd„„Zbe�rejcZbdÚd…d†„ZdejGd‡ƒjHZedˆd‰„ ZfdŠd‹„ ZgdŒd�„ Zhe�rHejiZhdÛdŽd�„Zjd�d‘„ Zkd’d“„ ZldÜd”d•„ZmdÝd–d—„Znd˜d™„ ZoG dšd›„ d›e@ƒZpdœZqG d�dž„ džepƒZrdŸd „ Zsd¡d¢„ Zte�r¾ejtZte*d£e+d¤e,d¥iZud¦d§„ Zvd¨d©„ Zwe�rêejwZwdªd«„ Zxd¬d­„ Zye�rejyZyd®d¯„ Zzd°d±„ Z{d²d³„ Z|e�r*ej|Z|d´dµ„ Z}e�r>ej}Z}d¶d·„ Z~d¸d¹„ Zdºd»„ Z€G d¼d½„ d½e@ƒZ�G d¾d¿„ d¿e@ƒZ‚e�jƒe�j„e‚jƒe‚j„iZ…dÀdÁ„ Z†dS )ÞzÕTools for creating `messages
<http://www.mongodb.org/display/DOCS/Mongo+Wire+Protocol>`_ to be sent to
MongoDB.

.. note:: This module is for internal use and is generally not needed by
   application developers.
é    N)ÚCodecOptionsÚdecodeÚencodeÚ_dict_to_bsonÚ_make_c_string)ÚDEFAULT_CODEC_OPTIONS)Ú_inflate_bsonÚDEFAULT_RAW_BSON_OPTIONS)ÚbÚStringIO)ÚSON)Ú	_cmessageTF)ÚConfigurationErrorÚCursorNotFoundÚDocumentTooLargeÚExecutionTimeoutÚInvalidOperationÚNotMasterErrorÚOperationFailureÚProtocolError)ÚDEFAULT_READ_CONCERN)ÚReadPreference)ÚWriteConcerniÿÿÿl        iþ?  é   é   ó    ó   ó    s     s       s           s       ÿÿÿÿs   documents     s   updates     s   deletes     Ú	documentsÚupdatesZdeletes)ÚinsertÚupdateÚdeletez%s.%sÚreplace)Zunicode_decode_error_handlerc               C   s   t jttƒS )z(Generate a pseudo random 32 bit integer.)ÚrandomÚrandintÚ	MIN_INT32Ú	MAX_INT32© r(   r(   ú2/tmp/pip-build-20mum3z4/pymongo/pymongo/message.pyÚ_randintZ   s    r*   c             C   sX   |j }|j}|j}|rT|tjj ks4|i gks4|dkrTd| krJtd| fgƒ} |j| d< | S )z-Add $readPreference to spec when appropriate.r   z$queryz$readPreferenceéÿÿÿÿ)ÚmodeÚtag_setsÚmax_stalenessr   ZSECONDARY_PREFERREDr   Údocument)ÚspecÚread_preferencer,   r-   r.   r(   r(   r)   Ú_maybe_add_read_preference_   s    

r2   c             C   s   t | ƒ| jjdœS )z<Convert an Exception into a failure document for publishing.)ÚerrmsgZerrtype)ÚstrÚ	__class__Ú__name__)Ú	exceptionr(   r(   r)   Ú_convert_exceptiont   s    r8   c       	      C   s  |j ddƒ}d|dœ}|j d|j ddƒƒ}|r„|j dƒrN|d	dd
idœ|d< n6d|j ddƒ|dœ}d|krv|d |d< |g|d< |S | dkržt|d ƒ|d< nv| dk�rd|krÆd|d dœg|d< nN|j dƒdkoÚ|dk�r|d d }|d j d|d j dƒƒ}d|dœg|d< |S )z7Convert a legacy write result to write commmand format.Únr   r   )Úokr9   r3   ÚerrÚ Zwtimeouté@   T)r3   ÚcodeÚerrInfoZwriteConcernErrorr>   é   )Úindexr>   r3   r?   ZwriteErrorsr    r   r!   Zupserted)rA   Ú_idZupdatedExistingFr   ÚurB   Úq)ÚgetÚlen)	Ú	operationÚcommandÚresultZaffectedÚresr3   Úerrorr!   rB   r(   r(   r)   Ú_convert_write_resultz   s2    




rL   ÚtailableÚoplogReplayr@   ÚnoCursorTimeouté   Ú	awaitDataé    ÚallowPartialResultsé€   ú$queryÚfilterú$orderbyÚsortú$hintÚhintú$commentÚcommentú$maxScanÚmaxScanú
$maxTimeMSÚ	maxTimeMSú$maxÚmaxú$minÚminú
$returnKeyÚ	returnKeyú$showRecordIdÚshowRecordIdú$showDiskLocú	$snapshotÚsnapshotc
                sì   t d| fgƒ}
d|krT|
jdd„ |jƒ D ƒƒ d|
kr@|
jdƒ d|
kr\|
jdƒ n||
d< |rh||
d< |rt||
d	< |r”t|ƒ|
d
< |dk r”d|
d< |r ||
d< |jr¼|	o®|	j r¼|j|
d< |rÈ||
d< ˆ rè|
j‡ fdd„tjƒ D ƒƒ |
S )z!Generate a find command document.Úfindz$queryc             S   s,   g | ]$\}}|t kr t | |fn||f‘qS r(   )Ú
_MODIFIERS)Ú.0ÚkeyÚvalr(   r(   r)   ú
<listcomp>½   s   z%_gen_find_command.<locals>.<listcomp>z$explainz$readPreferencerV   Ú
projectionÚskipÚlimitr   TZsingleBatchÚ	batchSizeÚreadConcernÚ	collationc                s    g | ]\}}ˆ |@ r|d f‘qS )Tr(   )rn   Úoptrp   )Úoptionsr(   r)   rq   Õ   s   )	r   r!   ÚitemsÚpopÚabsÚlevelÚin_transactionr/   Ú_OPTIONS)Úcollr0   rr   rs   rt   Ú
batch_sizery   Úread_concernrw   ÚsessionÚcmdr(   )ry   r)   Ú_gen_find_command¸   s6    


r…   c             C   s4   t d| fd|fgƒ}|r ||d< |dk	r0||d< |S )z$Generate a getMore command document.ÚgetMoreZ
collectionru   Nr`   )r   )Ú	cursor_idr€   r�   Úmax_await_time_msr„   r(   r(   r)   Ú_gen_get_more_commandÛ   s    r‰   c               @   sF   e Zd ZdZdZdZdZdd„ Zdd„ Zdd„ Z	dd„ Z
ddd„ZdS ) Ú_QueryzA query operation.ÚflagsÚdbr€   Úntoskipr0   ÚfieldsÚcodec_optionsr1   rt   r�   Únamer‚   rw   rƒ   ÚclientÚ_as_commandNc             C   sd   || _ || _|| _|| _|| _|| _|| _|| _|| _|	| _	|
| _
|| _|| _|| _d| _d | _d S )Nrl   )r‹   rŒ   r€   r�   r0   rŽ   r�   r1   r‚   rt   r�   rw   rƒ   r‘   r�   r’   )Úselfr‹   rŒ   r€   r�   r0   rŽ   r�   r1   rt   r�   r‚   rw   rƒ   r‘   r(   r(   r)   Ú__init__ò   s     z_Query.__init__c             C   s   t | j| jf S )N)Ú_UJOINrŒ   r€   )r“   r(   r(   r)   Ú	namespace  s    z_Query.namespacec             C   sn   d}|j dkr|s6d}n| jjs6td| jj|j f ƒ‚|j dk rZ| jd k	rZtd|j f ƒ‚|j| j| jƒ |S )NFé   TzDread concern level of %s is not valid with a max wire version of %d.é   zDSpecifying a collation is unsupported with a max wire version of %d.)	Úmax_wire_versionr‚   Zok_for_legacyr   r}   rw   Úvalidate_sessionr‘   rƒ   )r“   Ú	sock_infoÚexhaustZuse_find_cmdr(   r(   r)   Úuse_command	  s    
z_Query.use_commandc             C   sú   | j dk	r| j S d| jk}t| j| j| j| j| j| j| j| j	| j
| jƒ
}|r`d| _td|fgƒ}| j}|r¬|j|d| jƒ | r¬|jjr¬|jdk	r¬|j r¬|j|jdi ƒd< |j||| jƒ | j}|jrè|jj rè|jj| j|d| jƒ}|| jf| _ | j S )z.Return a find command document for this query.Nz$explainÚexplainFrv   ZafterClusterTime)r’   r0   r…   r€   rŽ   r�   rt   r�   r‹   r‚   rw   rƒ   r�   r   Ú	_apply_tor1   ry   Zcausal_consistencyZoperation_timer~   Ú
setdefaultÚsend_cluster_timer‘   Ú
_encrypterÚ_bypass_auto_encryptionÚencryptrŒ   r�   )r“   r›   rž   r„   rƒ   r‘   r(   r(   r)   Ú
as_command  s2    



z_Query.as_commandFc          
   C   sî   |r| j dB }n| j }| jƒ }| j}|r‚| j|ƒd }|jrntd|| j| j|d| j|j	d�\}}}	}
|||	fS t
| jdf }d	}n2| jdkr�dp”| j}| jr´|r®t| j|ƒ}n| j}|jrÆt|| jƒ}t||| j|||rÜdn| j| j|j	d�S )
z6Get a query message, possibly setting the slaveOk bit.r—   r   F)Úctxz$cmdr   r   Nr+   )r‹   r–   r0   r¥   Úop_msg_enabledÚ_op_msgrŒ   r1   r�   Úcompression_contextr•   r�   rt   rd   Z	is_mongosr2   Úqueryr�   rŽ   )r“   Úset_slave_okr›   Úuse_cmdr‹   Únsr0   Ú
request_idÚmsgÚsizeÚ_Ú	ntoreturnr(   r(   r)   Úget_messageA  s4    
z_Query.get_message)r‹   rŒ   r€   r�   r0   rŽ   r�   r1   rt   r�   r�   r‚   rw   rƒ   r‘   r’   )F)r6   Ú
__module__Ú__qualname__Ú__doc__Ú	__slots__Úexhaust_mgrr‡   r”   r–   r�   r¥   r³   r(   r(   r(   r)   rŠ   æ   s      #rŠ   c               @   sB   e Zd ZdZdZdZdd„ Zdd„ Zdd„ Zdd„ Z	ddd„Z
dS )Ú_GetMorezA getmore operation.rŒ   r€   r²   r‡   rˆ   r�   r1   rƒ   r‘   r¸   r’   r†   c             C   sF   || _ || _|| _|| _|| _|| _|| _|| _|	| _|
| _	d | _
d S )N)rŒ   r€   r²   r‡   r�   r1   rƒ   r‘   rˆ   r¸   r’   )r“   rŒ   r€   r²   r‡   r�   r1   rƒ   r‘   rˆ   r¸   r(   r(   r)   r”   s  s    z_GetMore.__init__c             C   s   t | j| jf S )N)r•   rŒ   r€   )r“   r(   r(   r)   r–   ‚  s    z_GetMore.namespacec             C   s    |j | j| jƒ |jdko| S )Nr—   )rš   r‘   rƒ   r™   )r“   r›   rœ   r(   r(   r)   r�   …  s    z_GetMore.use_commandc             C   sŽ   | j dk	r| j S t| j| j| j| jƒ}| jr>| jj|d| jƒ |j	|| j| j
ƒ | j
}|jr||jj r||jj| j|d| jƒ}|| jf| _ | j S )z1Return a getMore command document for this query.NF)r’   r‰   r‡   r€   r²   rˆ   rƒ   rŸ   r1   r¡   r‘   r¢   r£   r¤   rŒ   r�   )r“   r›   r„   r‘   r(   r(   r)   r¥   ‰  s    


z_GetMore.as_commandFc          
   C   s�   | j ƒ }|j}|r~| j|ƒd }|jrVtd|| jddd| j|jd�\}}}	}
|||	fS t| jdf }td|dd|d| j|d�S t	|| j
| j|ƒS )zGet a getmore message.r   NF)r¦   z$cmdr   r+   )r–   r©   r¥   r§   r¨   rŒ   r�   r•   rª   Úget_morer²   r‡   )r“   Zdummy0r›   r¬   r­   r¦   r0   r®   r¯   r°   r±   r(   r(   r)   r³   Ÿ  s    

z_GetMore.get_messageN)rŒ   r€   r²   r‡   rˆ   r�   r1   rƒ   r‘   r¸   r’   )F)r6   r´   rµ   r¶   r·   r�   r”   r–   r�   r¥   r³   r(   r(   r(   r)   r¹   j  s     r¹   c                   s*   e Zd Z‡ fdd„Zd‡ fdd„	Z‡  ZS )Ú_RawBatchQueryc                s   t t| ƒj||ƒ dS )NF)Úsuperr»   r�   )r“   Úsocket_inforœ   )r5   r(   r)   r�   µ  s    z_RawBatchQuery.use_commandFc                s   t t| ƒj||dƒS )NF)r¼   r»   r³   )r“   r«   r›   r¬   )r5   r(   r)   r³   »  s    
z_RawBatchQuery.get_message)F)r6   r´   rµ   r�   r³   Ú__classcell__r(   r(   )r5   r)   r»   ´  s   r»   c                   s&   e Zd Zdd„ Zd‡ fdd„	Z‡  ZS )Ú_RawBatchGetMorec             C   s   dS )NFr(   )r“   r½   rœ   r(   r(   r)   r�   Â  s    z_RawBatchGetMore.use_commandFc                s   t t| ƒj||dƒS )NF)r¼   r¿   r³   )r“   r«   r›   r¬   )r5   r(   r)   r³   Å  s    
z_RawBatchGetMore.get_message)F)r6   r´   rµ   r�   r³   r¾   r(   r(   )r5   r)   r¿   Á  s   r¿   c               @   s<   e Zd ZdZdd„ Zedd„ ƒZdd„ Zdd	„ Zd
d„ Z	dS )Ú_CursorAddresszEThe server address (host, port) of a cursor, with namespace property.c             C   s   t j| |ƒ}||_|S )N)ÚtupleÚ__new__Ú_CursorAddress__namespace)ÚclsÚaddressr–   r“   r(   r(   r)   rÂ   Î  s    z_CursorAddress.__new__c             C   s   | j S )zThe namespace this cursor.)rÃ   )r“   r(   r(   r)   r–   Ó  s    z_CursorAddress.namespacec             C   s   | | j f jƒ S )N)rÃ   Ú__hash__)r“   r(   r(   r)   rÆ   Ø  s    z_CursorAddress.__hash__c             C   s*   t |tƒr&t| ƒt|ƒko$| j|jkS tS )N)Ú
isinstancerÀ   rÁ   r–   ÚNotImplemented)r“   Úotherr(   r(   r)   Ú__eq__Ý  s    
z_CursorAddress.__eq__c             C   s
   | |k S )Nr(   )r“   rÉ   r(   r(   r)   Ú__ne__ã  s    z_CursorAddress.__ne__N)
r6   r´   rµ   r¶   rÂ   Úpropertyr–   rÆ   rÊ   rË   r(   r(   r(   r)   rÀ   Ë  s   rÀ   z<iiiiiiBé   c             C   s>   |j |ƒ}tƒ }ttt|ƒ |dd| t|ƒ|jƒ}||| fS )zDTakes message data, compresses it, and adds an OP_COMPRESSED header.r   iÜ  )Úcompressr*   Ú_pack_compression_headerÚ_COMPRESSION_HEADER_SIZErF   Zcompressor_id)rG   Údatar¦   Ú
compressedr®   Úheaderr(   r(   r)   Ú	_compressê  s    

rÔ   c             C   s<   t dgƒ}|j|ƒ | jddƒ}td|d d dd|dtƒS )	z$Data to send to do a lastError.
    Úgetlasterrorr   Ú.r   z.$cmdN)rÕ   r   r+   )r   r!   Úsplitrª   r   )r–   Úargsr„   Zsplitnsr(   r(   r)   Ú__last_errorú  s
    

rÙ   z<iiiic             C   s(   t ƒ }tdt|ƒ |d| ƒ}||| fS )ztTakes message data and adds a message header based on the operation.

    Returns the resultant message string.
    rP   r   )r*   Ú_pack_headerrF   )rG   rÑ   ÚridÚmessager(   r(   r)   Ú__pack_message  s    rÝ   z<ic                sŠ   t ‰t|ƒdkr<ˆ|d ˆ ˆƒ}djdt| ƒ|gƒt|ƒfS ‡ ‡‡fdd„|D ƒ}|s^tdƒ‚djt|ƒt| ƒdj|ƒgƒttt|ƒƒfS )zGet an OP_INSERT messager   r   r   s       c                s   g | ]}ˆ|ˆ ˆƒ‘qS r(   r(   )rn   Údoc)Ú
check_keysr   Úoptsr(   r)   rq     s    z_insert.<locals>.<listcomp>zcannot do an empty bulk insert)r   rF   Újoinr   r   Ú	_pack_intrb   Úmap)Úcollection_nameÚdocsrß   r‹   rà   Úencodedr(   )rß   r   rà   r)   Ú_insert  s    rç   c       
      C   s.   t | ||||ƒ\}}td||ƒ\}}	||	|fS )z9Internal compressed unacknowledged insert message helper.iÒ  )rç   rÔ   )
rä   rå   rß   Úcontinue_on_errorrà   r¦   Ú	op_insertÚmax_bson_sizerÛ   r¯   r(   r(   r)   Ú_insert_compressed'  s    rë   c             C   sN   t | ||||ƒ\}}td|ƒ\}	}
|rDt| |ƒ\}	}}|	|
| |fS |	|
|fS )zInternal insert message helper.iÒ  )rç   rÝ   rÙ   )rä   rå   rß   ÚsafeÚlast_error_argsrè   rà   ré   rê   rÛ   r¯   Úgler±   r(   r(   r)   Ú_insert_uncompressed0  s    rï   c             C   s*   |rt | |||||ƒS t| ||||||ƒS )zGet an **insert** message.)rë   rï   )rä   rå   rß   rì   rí   rè   rà   r¦   r(   r(   r)   r    >  s
    
r    c       
      C   sX   d}|r|d7 }|r|d7 }t }||||ƒ}	djtt| ƒt|ƒ||d|ƒ|	gƒt|	ƒfS )zGet an OP_UPDATE message.r   r   r   r   F)r   rá   Ú_ZERO_32r   râ   rF   )
rä   ÚupsertÚmultir0   rÞ   rß   rà   r‹   r   Zencoded_updater(   r(   r)   Ú_updateH  s    
ró   c             C   s2   t | ||||||ƒ\}}	td||ƒ\}
}|
||	fS )z9Internal compressed unacknowledged update message helper.iÑ  )ró   rÔ   )rä   rñ   rò   r0   rÞ   rß   rà   r¦   Ú	op_updaterê   rÛ   r¯   r(   r(   r)   Ú_update_compressedY  s    rõ   c	             C   sR   t | ||||||ƒ\}	}
td|	ƒ\}}|rHt| |ƒ\}}}||| |
fS |||
fS )zInternal update message helper.iÑ  )ró   rÝ   rÙ   )rä   rñ   rò   r0   rÞ   rì   rí   rß   rà   rô   rê   rÛ   r¯   rî   r±   r(   r(   r)   Ú_update_uncompressedb  s    rö   c
       
   
   C   s2   |	rt | |||||||	ƒS t| ||||||||ƒ	S )zGet an **update** message.)rõ   rö   )
rä   rñ   rò   r0   rÞ   rì   rí   rß   rà   r¦   r(   r(   r)   r!   p  s
    
r!   z<IBz<Bc                s¶   t |dˆƒ}t| dƒ}t|ƒ}d}	|ržtdƒ}
t|ƒ}‡ ‡fdd„|D ƒ}t|ƒtdd„ |D ƒƒ d }t|ƒ}||7 }td	d„ |D ƒƒ}	|||
||g| }n||g}d
j|ƒ||	fS )zäGet a OP_MSG message.

    Note: this method handles multiple documents in a type one payload but
    it does not perform batch splitting and the total message size is
    only checked *after* generating the entire message.
    Fr   r   c                s   g | ]}t |ˆ ˆƒ‘qS r(   )r   )rn   rÞ   )rß   rà   r(   r)   rq   �  s    z%_op_msg_no_header.<locals>.<listcomp>c             s   s   | ]}t |ƒV  qd S )N)rF   )rn   rÞ   r(   r(   r)   ú	<genexpr>Ž  s    z$_op_msg_no_header.<locals>.<genexpr>r—   c             s   s   | ]}t |ƒV  qd S )N)rF   )rn   rÞ   r(   r(   r)   r÷   ‘  s    r   )	r   Ú_pack_op_msg_flags_typerF   Ú
_pack_byter   Úsumrâ   rb   rá   )r‹   rH   Ú
identifierrå   rß   rà   ræ   Z
flags_typeÚ
total_sizeÚmax_doc_sizeZtype_oneZcstringZencoded_docsr°   Zencoded_sizerÑ   r(   )rß   rà   r)   Ú_op_msg_no_header~  s     
rþ   c             C   s4   t | |||||ƒ\}}}	td||ƒ\}
}|
|||	fS )zInternal OP_MSG message helper.iÝ  )rþ   rÔ   )r‹   rH   rû   rå   rß   rà   r¦   r¯   rü   rê   rÛ   r(   r(   r)   Ú_op_msg_compressed™  s    rÿ   c             C   s2   t | |||||ƒ\}}}td|ƒ\}	}
|	|
||fS )z*Internal compressed OP_MSG message helper.iÝ  )rþ   rÝ   )r‹   rH   rû   rå   rß   rà   rÑ   rü   rê   r®   Z
op_messager(   r(   r)   Ú_op_msg_uncompressed¢  s    r   c             C   s¼   ||d< |dk	r<d|kr<|r2|j  r2tjj|d< n
|j|d< tt|ƒƒ}ytj|ƒ}	|j|	ƒ}
W n t	k
r|   d}	d}
Y nX z*|r˜t
| ||	|
|||ƒS t| ||	|
||ƒS |	r¶|
||	< X dS )zGet a OP_MSG message.z$dbNz$readPreferencer<   )r,   r   ZPRIMARY_PREFERREDr/   ÚnextÚiterÚ
_FIELD_MAPrE   r{   ÚKeyErrorrÿ   r   )r‹   rH   Zdbnamer1   Úslave_okrß   rà   r¦   r�   rû   rå   r(   r(   r)   r¨   ¬  s(    


r¨   c             C   s^   t |||ƒ}|rt |d|ƒ}	nd}	tt|ƒt|	ƒƒ}
djt| ƒt|ƒt|ƒt|ƒ||	gƒ|
fS )zGet an OP_QUERY message.Fr   )r   rb   rF   rá   râ   r   )ry   rä   Únum_to_skipÚnum_to_returnrª   Úfield_selectorrà   rß   ræ   Zefsrê   r(   r(   r)   Ú_queryÊ  s    r	  c	          	   C   s4   t | |||||||ƒ\}	}
td|	|ƒ\}}|||
fS )z)Internal compressed query message helper.iÔ  )r	  rÔ   )ry   rä   r  r  rª   r  rà   rß   r¦   Úop_queryrê   rÛ   r¯   r(   r(   r)   Ú_query_compressedÜ  s    
r  c          	   C   s2   t | |||||||ƒ\}}	td|ƒ\}
}|
||	fS )zInternal query message helper.iÔ  )r	  rÝ   )ry   rä   r  r  rª   r  rà   rß   r
  rê   rÛ   r¯   r(   r(   r)   Ú_query_uncompressedí  s    
r  c	       	   
   C   s2   |rt | ||||||||ƒ	S t| |||||||ƒS )zGet a **query** message.)r  r  )	ry   rä   r  r  rª   r  rà   rß   r¦   r(   r(   r)   rª   ÿ  s    
rª   z<qc             C   s   dj tt| ƒt|ƒt|ƒgƒS )zGet an OP_GET_MORE message.r   )rá   rð   r   râ   Ú_pack_long_long)rä   r  r‡   r(   r(   r)   Ú	_get_more  s
    r  c             C   s   t dt| ||ƒ|ƒS )z+Internal compressed getMore message helper.iÕ  )rÔ   r  )rä   r  r‡   r¦   r(   r(   r)   Ú_get_more_compressed  s    r  c             C   s   t dt| ||ƒƒS )z Internal getMore message helper.iÕ  )rÝ   r  )rä   r  r‡   r(   r(   r)   Ú_get_more_uncompressed  s    r  c             C   s   |rt | |||ƒS t| ||ƒS )zGet a **getMore** message.)r  r  )rä   r  r‡   r¦   r(   r(   r)   rº   %  s    rº   c             C   s.   t |d|ƒ}djtt| ƒt|ƒ|gƒt|ƒfS )zGet an OP_DELETE message.Fr   )r   rá   rð   r   râ   rF   )rä   r0   rà   r‹   ræ   r(   r(   r)   Ú_delete-  s    r  c       	      C   s,   t | |||ƒ\}}td||ƒ\}}|||fS )z9Internal compressed unacknowledged delete message helper.iÖ  )r  rÔ   )	rä   r0   rà   r‹   r¦   Ú	op_deleterê   rÛ   r¯   r(   r(   r)   Ú_delete_compressed7  s    r  c             C   sL   t | |||ƒ\}}td|ƒ\}}	|rBt| |ƒ\}}
}||	|
 |fS ||	|fS )zInternal delete message helper.iÖ  )r  rÝ   rÙ   )rä   r0   rì   rí   rà   r‹   r  rê   rÛ   r¯   rî   r±   r(   r(   r)   Ú_delete_uncompressed>  s    r  c             C   s&   |rt | ||||ƒS t| |||||ƒS )zàGet a **delete** message.

    `opts` is a CodecOptions. `flags` is a bit vector that may contain
    the SingleRemove flag or not:

    http://docs.mongodb.org/meta-driver/latest/legacy/mongodb-wire-protocol/#op-delete
    )r  r  )rä   r0   rì   rí   rà   r‹   r¦   r(   r(   r)   r"   I  s    	r"   c             C   s6   t | ƒ}tjdd|  ƒj}|d|f| žŽ }td|ƒS )z#Get a **killCursors** message.
    z<iirD   r   i×  )rF   ÚstructÚStructÚpackrÝ   )Z
cursor_idsZnum_cursorsr  Zop_kill_cursorsr(   r(   r)   Úkill_cursorsX  s    r  c               @   s    e Zd ZdZd.Zdd„ Zdd„ Zdd„ Zdd„ Ze	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d%d&„ Zd'd(„ Zd)d*„ Zd+d,„ Zd-S )/Ú_BulkWriteContextzCA wrapper around SocketInfo for use with write splitting functions.Údb_namerH   r›   Úop_idr�   ÚfieldÚpublishÚ
start_timeÚ	listenersrƒ   rÎ   Úop_typeÚcodecc	       	      C   s|   || _ || _|| _|| _|| _|j| _tt|ƒƒ| _	t
| j	 | _| jrPtjjƒ nd | _|| _|jrfdnd| _|| _|| _d S )NTF)r  rH   r›   r  r  Úenabled_for_commandsr  r  r  r�   r  r  ÚdatetimeÚnowr  rƒ   r©   rÎ   r   r!  )	r“   Zdatabase_namerH   r›   Zoperation_idr  rƒ   r   r!  r(   r(   r)   r”   h  s    z_BulkWriteContext.__init__c             C   sB   | j d }t|| j| j|| j| j| ƒ\}}}|s8tdƒ‚|||fS )Nz.$cmdzcannot do an empty bulk write)r  Ú_do_bulk_write_commandr   rH   rß   r!  r   )r“   rå   r–   r®   r¯   Úto_sendr(   r(   r)   Ú_batch_commandx  s    
z _BulkWriteContext._batch_commandc             C   s4   | j |ƒ\}}}| j|||ƒ}|j|| jƒ ||fS )N)r'  Úwrite_commandZ_process_responserƒ   )r“   rå   r‘   r®   r¯   r&  rI   r(   r(   r)   Úexecute�  s    z_BulkWriteContext.executec             C   s&   | j |ƒ\}}}| j||dd|ƒ |S )Nr   F)r'  Úlegacy_write)r“   rå   r‘   r®   r¯   r&  r(   r(   r)   Úexecute_unack‡  s    z_BulkWriteContext.execute_unackc             C   s
   | j tkS )z-Should we check keys for this operation type?)r   Ú_INSERT)r“   r(   r(   r)   rß   ‘  s    z_BulkWriteContext.check_keysc             C   s   | j jS )z#A proxy for SockInfo.max_bson_size.)r›   rê   )r“   r(   r(   r)   rê   –  s    z_BulkWriteContext.max_bson_sizec             C   s   | j r| jjd S | jjS )z&A proxy for SockInfo.max_message_size.rP   )rÎ   r›   Úmax_message_size)r“   r(   r(   r)   r-  ›  s    z"_BulkWriteContext.max_message_sizec             C   s   | j jS )z*A proxy for SockInfo.max_write_batch_size.)r›   Úmax_write_batch_size)r“   r(   r(   r)   r.  £  s    z&_BulkWriteContext.max_write_batch_sizec             C   s   | j S )z:The maximum size of a BSON command before batch splitting.)rê   )r“   r(   r(   r)   Úmax_split_size¨  s    z _BulkWriteContext.max_split_sizec             C   s*   |rt d|| jjƒ\}}| j|||||ƒS )NiÒ  )rÔ   r›   r©   r*  )r“   r®   r¯   rý   Úacknowledgedrå   rÎ   r(   r(   r)   Úlegacy_bulk_insert­  s
    z$_BulkWriteContext.legacy_bulk_insertc             C   sø   | j r,tjjƒ | j }| j||ƒ}tjjƒ }z¸y\| jj||||ƒ}	| j rˆtjjƒ | | }|	dk	rrt| j||	ƒ}
nddi}
| j	||
|ƒ W nV t
k
rà } z:| j rÎtjjƒ | | }| j|t| j||jƒ|ƒ ‚ W Y dd}~X nX W dtjjƒ | _X |	S )zKA proxy for SocketInfo.legacy_write that handles event publishing.
        Nr:   r   )r  r#  r$  r  Ú_startr›   r*  rL   r�   Ú_succeedr   Ú_failÚdetails)r“   r®   r¯   rý   r0  rå   Údurationr„   ÚstartrI   ÚreplyÚexcr(   r(   r)   r*  µ  s0    
z_BulkWriteContext.legacy_writec             C   sÊ   | j r,tjjƒ | j }| j||ƒ tjjƒ }zŠy8| jj||ƒ}| j rdtjjƒ | | }| j|||ƒ W nL tk
r² } z0| j r tjjƒ | | }| j	||j
|ƒ ‚ W Y dd}~X nX W dtjjƒ | _X |S )zLA proxy for SocketInfo.write_command that handles event publishing.
        N)r  r#  r$  r  r2  r›   r(  r3  r   r4  r5  )r“   r®   r¯   rå   r6  r7  r8  r9  r(   r(   r)   r(  Ô  s     
z_BulkWriteContext.write_commandc             C   s4   | j jƒ }||| j< | jj|| j|| jj| jƒ |S )zPublish a CommandStartedEvent.)	rH   Úcopyr  r  Úpublish_command_startr  r›   rÅ   r  )r“   r®   rå   r„   r(   r(   r)   r2  é  s    

z_BulkWriteContext._startc             C   s"   | j j||| j|| jj| jƒ dS )z Publish a CommandSucceededEvent.N)r  Úpublish_command_successr�   r›   rÅ   r  )r“   r®   r8  r6  r(   r(   r)   r3  ò  s    z_BulkWriteContext._succeedc             C   s"   | j j||| j|| jj| jƒ dS )zPublish a CommandFailedEvent.N)r  Úpublish_command_failurer�   r›   rÅ   r  )r“   r®   Úfailurer6  r(   r(   r)   r4  ø  s    z_BulkWriteContext._failN)r  rH   r›   r  r�   r  r  r  r  rƒ   rÎ   r   r!  )r6   r´   rµ   r¶   r·   r”   r'  r)  r+  rÌ   rß   rê   r-  r.  r/  r1  r*  r(  r2  r3  r4  r(   r(   r(   r)   r  a  s&     	
	r  i    c               @   s4   e Zd Zf Zdd„ Zdd„ Zdd„ Zedd„ ƒZd	S )
Ú_EncryptedBulkWriteContextc             C   sd   | j d }t|| j| j|| j| j| ƒ\}}|s6tdƒ‚|jddƒd }tt	|ƒ|d … t
ƒ}||fS )Nz.$cmdzcannot do an empty bulk writer   r—   é	   )r  Ú_encode_batched_write_commandr   rH   rß   r!  r   rA   r   Ú
memoryviewr	   )r“   rå   r–   r¯   r&  Z	cmd_startr„   r(   r(   r)   r'  	  s    
z)_EncryptedBulkWriteContext._batch_commandc             C   s0   | j |ƒ\}}| jj| j|t| j|d�}||fS )N)r�   rƒ   r‘   )r'  r›   rH   r  Ú_UNICODE_REPLACE_CODEC_OPTIONSrƒ   )r“   rå   r‘   r„   r&  rI   r(   r(   r)   r)    s
    z"_EncryptedBulkWriteContext.executec             C   s2   | j |ƒ\}}| jj| j|tdd�| j|d� |S )Nr   )Úw)Zwrite_concernrƒ   r‘   )r'  r›   rH   r  r   rƒ   )r“   rå   r‘   r„   r&  r(   r(   r)   r+    s
    z(_EncryptedBulkWriteContext.execute_unackc             C   s   t S )z Reduce the batch splitting size.)Ú_MAX_SPLIT_SIZE_ENC)r“   r(   r(   r)   r/  %  s    z)_EncryptedBulkWriteContext.max_split_sizeN)	r6   r´   rµ   r·   r'  r)  r+  rÌ   r/  r(   r(   r(   r)   r?    s
   r?  c             C   s,   | dkrt d||f ƒ‚nt d| f ƒ‚dS )z-Internal helper for raising DocumentTooLarge.r    zfBSON document too large (%d bytes) - the connected server supports BSON document sizes up to %d bytes.z%r command document too largeN)r   )rG   Zdoc_sizeÚmax_sizer(   r(   r)   Ú_raise_document_too_large+  s    rG  c                sì  ‡ ‡fdd„}|p| }	d}
t ƒ }|jtjdt|ƒƒƒ |jtˆ ƒƒ |jƒ  }}d}g }t}|jol|pj|	 }�x|D �]}||||ƒ}t	|ƒ}||j
k}||7 }||jk rÌ| rÌ|j|ƒ |j|ƒ d}qv|�rNy>|rèd|jƒ  }}n||jƒ |	ƒ\}}|j||d|	||ƒ W n< tk
�rL } z|�r0|}
n|�s:dS ‚ W Y dd}~X nX |�rbtd||j
ƒ || }|j|ƒ |jƒ  |j|ƒ |g}qvW |�sžtd	ƒ‚|�r´d|jƒ  }}n||jƒ |ƒ\}}|j||d|||ƒ |
dk	�rè|
‚dS )
z*Insert `docs` using multiple batches.
    c                s2   t d| ƒ\}}|r*tˆ ˆƒ\}}}||7 }||fS )z6Build the insert message with header and GLE.
        iÒ  )rÝ   rÙ   )Zinsert_messageÚ	send_safer®   Zfinal_messageÚerror_messager±   )rä   rí   r(   r)   Ú_insert_message=  s    z+_do_batched_insert.<locals>._insert_messageNz<iFTr   r    zcannot do an empty bulk insert)r   Úwriter  r  Úintr   Útellr   rÎ   rF   rê   r-  ÚappendÚgetvaluer1  r   rG  ÚseekÚtruncater   )rä   rå   rß   rì   rí   rè   rà   r¦   rJ  rH  Z
last_errorrÑ   Zmessage_lengthZ	begin_locZhas_docsr&  r   rÎ   rÞ   ræ   Zencoded_lengthZ	too_largerÛ   r¯   r9  r®   r(   )rä   rí   r)   Ú_do_batched_insert8  sd    








rR  s
   documents s   updates s   deletes c             C   s|  |j }|j}	|j}
|rdnd}|j|ƒ |jdƒ |jt|d|ƒƒ |jdƒ |jƒ }|jdƒ y|jt|  ƒ W n tk
rŽ   tdƒ‚Y nX | t	t
fkr d}g }d}x¦|D ]ž}t|||ƒ}t|ƒ}|jƒ | }|dkoà||
k}| oî||k}|sú|�rttjƒ ƒ|  }t|t|ƒ|ƒ ||
k�r&P |j|ƒ |j|ƒ |d7 }||	kr®P q®W |jƒ }|j|ƒ |jt|| ƒƒ ||fS )	zCreate a batched OP_MSG write.s       s      r   Fó   zUnknown commandr   r   )rê   r.  r-  rK  r   rM  Ú_OP_MSG_MAPr  r   Ú_UPDATEÚ_DELETErF   Úlistr  ÚkeysrG  rN  rP  râ   )rG   rH   rå   rß   Úackrà   r¦   Úbufrê   r.  r-  r‹   Zsize_locationr&  ÚidxrÞ   ÚvalueZ
doc_lengthZnew_message_sizeÚdoc_too_largeZunacked_doc_too_largeÚwrite_opÚlengthr(   r(   r)   Ú_batched_op_msg_impl—  sN    









r`  c       
   	   C   s,   t ƒ }t| |||||||ƒ\}}	|jƒ |fS )zOEncode the next batched insert, update, or delete operation
    as OP_MSG.
    )r   r`  rO  )
rG   rH   rå   rß   rY  rà   r¦   rZ  r&  r±   r(   r(   r)   Ú_encode_batched_op_msgØ  s    ra  c             C   s6   t | ||||||ƒ\}}td||jjƒ\}	}
|	|
|fS )z]Create the next batched insert, update, or delete operation
    with OP_MSG, compressed.
    iÝ  )ra  rÔ   r›   r©   )rG   rH   rå   rß   rY  rà   r¦   rÑ   r&  r®   r¯   r(   r(   r)   Ú_batched_op_msg_compressedæ  s    rb  c          	   C   sx   t ƒ }|jtƒ |jdƒ t| |||||||ƒ\}}	|jdƒ tƒ }
|jt|
ƒƒ |jdƒ |jt|	ƒƒ |
|jƒ |fS )z"OP_MSG implementation entry point.s       Ý  r—   r   )r   rK  Ú_ZERO_64r`  rP  r*   râ   rO  )rG   rH   rå   rß   rY  rà   r¦   rZ  r&  r_  r®   r(   r(   r)   Ú_batched_op_msgõ  s    



rd  c             C   sf   | j ddƒd |d< d|kr2t|d jddƒƒ}nd}|jjrRt|||||||ƒS t|||||||ƒS )zRCreate the next batched insert, update, or delete operation
    using OP_MSG.
    rÖ   r   r   z$dbZwriteConcernrD  T)r×   ÚboolrE   r›   r©   rb  rd  )r–   rG   rH   rå   rß   rà   r¦   rY  r(   r(   r)   Ú_do_batched_op_msg  s    rf  c             C   s6   t | ||||||ƒ\}}td||jjƒ\}	}
|	|
|fS )zKCreate the next batched insert, update, or delete command, compressed.
    iÔ  )rA  rÔ   r›   r©   )r–   rG   rH   rå   rß   rà   r¦   rÑ   r&  r®   r¯   r(   r(   r)   Ú!_batched_write_command_compressed"  s    rg  c       
   	   C   s,   t ƒ }t| |||||||ƒ\}}	|jƒ |fS )z?Encode the next batched insert, update, or delete command.
    )r   Ú_batched_write_command_implrO  )
r–   rG   rH   rå   rß   rà   r¦   rZ  r&  r±   r(   r(   r)   rA  0  s    rA  c          	   C   sx   t ƒ }|jtƒ |jdƒ t| |||||||ƒ\}}	|jdƒ tƒ }
|jt|
ƒƒ |jdƒ |jt|	ƒƒ |
|jƒ |fS )z?Create the next batched insert, update, or delete command.
    s       Ô  r—   r   )r   rK  rc  rh  rP  r*   râ   rO  )r–   rG   rH   rå   rß   rà   r¦   rZ  r&  r_  r®   r(   r(   r)   Ú_batched_write_command=  s    



ri  c             C   s0   |j jrt| ||||||ƒS t| ||||||ƒS )z#Batched write commands entry point.)r›   r©   rg  ri  )r–   rG   rH   rå   rß   rà   r¦   r(   r(   r)   Ú_do_batched_write_commandX  s
    rj  c             C   s4   |j jdkr t| ||||||ƒS t| ||||||ƒS )z Bulk write commands entry point.r˜   )r›   r™   rf  rj  )r–   rG   rH   rå   rß   rà   r¦   r(   r(   r)   r%  b  s
    r%  c             C   sè  |j }|j}	|t }
|j}|jtƒ |jt| ƒƒ |jtƒ |jtƒ |j	ƒ }|jt
|ƒƒ |jddƒ |jƒ  y|jt| ƒ W n tk
rž   tdƒ‚Y nX |ttfkr°d}|j	ƒ d }g }d}xÌ|D ]Ä}tt|ƒƒ}t
|||ƒ}t|ƒ|
k}|�rttjƒ ƒ| }t|t|ƒ|ƒ |dk�o<|j	ƒ t|ƒ t|ƒ |k}||	k}|�sR|�rTP |jtƒ |j|ƒ |jtƒ |j|ƒ |j|ƒ |d7 }qÊW |jtƒ |j	ƒ }|j|ƒ |jt|| d ƒƒ |j|ƒ |jt|| ƒƒ ||fS )z(Create a batched OP_QUERY write command.r   r   zUnknown commandFr—   r   r+   )rê   r.  Ú_COMMAND_OVERHEADr/  rK  rð   r
   Ú_ZERO_8Ú_SKIPLIMrM  r   rP  rQ  Ú_OP_MAPr  r   rU  rV  r4   rF   rW  r  rX  rG  Ú_BSONOBJrN  Ú_ZERO_16râ   )r–   rG   rH   rå   rß   rà   r¦   rZ  rê   r.  Zmax_cmd_sizer/  Zcommand_startZ
list_startr&  r[  rÞ   ro   r\  r]  r^  Zenough_dataZenough_documentsr_  r(   r(   r)   rh  l  s^    












rh  c               @   sd   e Zd ZdZdZejdƒjZdZ	dd	„ Z
ddd„Zd
ed
dfdd„Zdd„ Zdd„ Zedd„ ƒZd
S )Ú_OpReplyz$A MongoDB OP_REPLY response message.r‹   r‡   Únumber_returnedr   z<iqiir   c             C   s   || _ || _|| _|| _d S )N)r‹   r‡   rr  r   )r“   r‹   r‡   rr  r   r(   r(   r)   r”   ¿  s    z_OpReply.__init__Nc             C   sÌ   | j d@ r>|dkrtdƒ‚d|f }d|ddœ}t|d|ƒ‚n†| j d@ rÄtj| jƒjƒ }|jd	dƒ |d
 jdƒr‚t	|d
 |ƒ‚n&|j
dƒdkr¨t|j
d
ƒ|j
dƒ|ƒ‚td|j
d
ƒ |j
dƒ|ƒ‚| jgS )aº  Check the response header from the database, without decoding BSON.

        Check the response for errors and unpack.

        Can raise CursorNotFound, NotMasterError, ExecutionTimeout, or
        OperationFailure.

        :Parameters:
          - `cursor_id` (optional): cursor_id we sent to get this response -
            used for raising an informative exception when we get cursor id not
            valid at server response.
        r   Nz"No cursor id for getMore operationzCursor not found, cursor id: %dr   é+   )r:   r3   r>   r   r:   z$errz
not masterr>   é2   zdatabase error: %s)r‹   r   r   ÚbsonZBSONr   r   r    Ú
startswithr   rE   r   r   )r“   r‡   r¯   ZerrobjZerror_objectr(   r(   r)   Úraw_responseÅ  s(    




z_OpReply.raw_responseFc             C   s,   | j |ƒ |rtj| j|ƒS tj| j||ƒS )ad  Unpack a response from the database and decode the BSON document(s).

        Check the response for errors and unpack, returning a dictionary
        containing the response data.

        Can raise CursorNotFound, NotMasterError, ExecutionTimeout, or
        OperationFailure.

        :Parameters:
          - `cursor_id` (optional): cursor_id we sent to get this response -
            used for raising an informative exception when we get cursor id not
            valid at server response
          - `codec_options` (optional): an instance of
            :class:`~bson.codec_options.CodecOptions`
        )rw  ru  Z
decode_allr   Ú_decode_all_selective)r“   r‡   r�   Úuser_fieldsÚlegacy_responser(   r(   r)   Úunpack_responseì  s
    
z_OpReply.unpack_responsec             C   s   | j ƒ }| jdkst‚|d S )zUnpack a command response.r   r   )r{  rr  ÚAssertionError)r“   rå   r(   r(   r)   Úcommand_response  s    z_OpReply.command_responsec             C   s   t ‚dS )z)Return the bytes of the command response.N)ÚNotImplementedError)r“   r(   r(   r)   Úraw_command_response
  s    z_OpReply.raw_command_responsec             C   s0   | j |ƒ\}}}}t|dd… ƒ}| ||||ƒS )z%Construct an _OpReply from raw bytes.é   N)ÚUNPACK_FROMÚbytes)rÄ   r¯   r‹   r‡   r±   rr  r   r(   r(   r)   Úunpack  s    z_OpReply.unpack)r‹   r‡   rr  r   )N)r6   r´   rµ   r¶   r·   r  r  Úunpack_fromr�  ÚOP_CODEr”   rw  rC  r{  r}  r  Úclassmethodrƒ  r(   r(   r(   r)   rq  ·  s   
'rq  c               @   sd   e Zd ZdZdZejdƒjZdZ	dd	„ Z
ddd„Zd
ed
dfdd„Zdd„ Zdd„ Zedd„ ƒZd
S )Ú_OpMsgz"A MongoDB OP_MSG response message.r‹   r‡   rr  Úpayload_documentz<IBiiÝ  c             C   s   || _ || _d S )N)r‹   rˆ  )r“   r‹   rˆ  r(   r(   r)   r”   #  s    z_OpMsg.__init__Nc             C   s   t ‚d S )N)r~  )r“   r‡   r(   r(   r)   rw  '  s    z_OpMsg.raw_responseFc             C   s   | s
t ‚tj| j||ƒS )zûUnpack a OP_MSG command response.

        :Parameters:
          - `cursor_id` (optional): Ignored, for compatibility with _OpReply.
          - `codec_options` (optional): an instance of
            :class:`~bson.codec_options.CodecOptions`
        )r|  ru  rx  rˆ  )r“   r‡   r�   ry  rz  r(   r(   r)   r{  *  s    
z_OpMsg.unpack_responsec             C   s   | j ƒ d S )zUnpack a command response.r   )r{  )r“   r(   r(   r)   r}  9  s    z_OpMsg.command_responsec             C   s   | j S )z)Return the bytes of the command response.)rˆ  )r“   r(   r(   r)   r  =  s    z_OpMsg.raw_command_responsec             C   sn   | j |ƒ\}}}|dkr&td|f ƒ‚|dkr<td|f ƒ‚t|ƒ|d krTtdƒ‚t|dd… ƒ}| ||ƒS )z#Construct an _OpMsg from raw bytes.r   zUnsupported OP_MSG flags (%r)z$Unsupported OP_MSG payload type (%r)r˜   z$Unsupported OP_MSG reply: >1 sectionN)r�  r   rF   r‚  )rÄ   r¯   r‹   Zfirst_payload_typeZfirst_payload_sizerˆ  r(   r(   r)   rƒ  A  s    z_OpMsg.unpack)r‹   r‡   rr  rˆ  )N)r6   r´   rµ   r¶   r·   r  r  r„  r�  r…  r”   rw  rC  r{  r}  r  r†  rƒ  r(   r(   r(   r)   r‡    s   
r‡  c
             C   sŒ  t d||d|d|||dtdddƒ}tt|ƒƒ}
|	j}|rBtjjƒ }|j|| ƒ\}}}|r‚tjjƒ | }|	j|||| j	ƒ tjjƒ }| j
||ƒ | j|ƒ}y|jd|ƒ}W np tk
�r } zR|�rtjjƒ | | }t|ttfƒrê|j}nt|ƒ}|	j|||
|| j	ƒ ‚ W Y dd}~X nX d|k�rB||jd||f dœddœ}n|�rP|d ni }d|d< |�rˆtjjƒ | | }|	j|||
|| j	ƒ |S )	zESimple query helper for retrieving a first (and possibly only) batch.r   NÚcursorz%s.%s)Z
firstBatchÚidr­   g      ð?)r‰  r:   r:   )rŠ   r   r  r  r"  r#  r$  r³   r;  rÅ   Úsend_messageZreceive_messager{  Ú	ExceptionrÇ   r   r   r5  r8   r=  r‡   r<  )r›   rŒ   r€   rª   r²   r  r�   r1   r„   r  r�   r  r7  r®   r¯   rý   Zencoding_durationr8  rå   r9  r6  r>  rI   r(   r(   r)   Ú_first_batchZ  sN    




r�  i   €)rM   r   )rN   r@   )rO   rP   )rQ   rR   )rS   rT   )rU   rV   )rW   rX   )rY   rZ   )r[   r\   )r]   r^   )r_   r`   )ra   rb   )rc   rd   )re   rf   )rg   rh   )ri   rh   )rj   rk   )NN)N)N)N)FN)F)FN)N)r   )r   N)‡r¶   r#  r$   r  ru  r   r   r   r   r   Zbson.codec_optionsr   Zbson.raw_bsonr   r	   Zbson.py3compatr
   r   Zbson.sonr   Zpymongor   Z_use_cÚImportErrorZpymongo.errorsr   r   r   r   r   r   r   r   Zpymongo.read_concernr   Zpymongo.read_preferencesr   Zpymongo.write_concernr   r'   r&   rk  r,  rU  rV  Z_EMPTYro  rl  rp  rð   rc  rm  rn  r  r•   rC  r*   r2   r8   rL   r   rm   r…   r‰   ÚobjectrŠ   r¹   r»   r¿   rÁ   rÀ   r  r  rÏ   rÐ   rÔ   rÙ   rÚ   rÝ   râ   rç   rë   rï   rJ  r    ró   rõ   rö   Z_update_messager!   rø   rù   rþ   rÿ   r   r¨   r	  r  r  Z_query_messagerª   r  r  r  r  Z_get_more_messagerº   r  r  r  r"   r  r  rE  r?  rG  rR  rT  r`  ra  rb  rd  rf  rg  rA  ri  rj  r%  rh  rq  r‡  r…  rƒ  Z_UNPACK_REPLYr�  r(   r(   r(   r)   Ú<module>   s.  
('
" J


	

		

		



	



	 #%RA
	

Kd: