3
]ð]\"  ã               @   s„   d 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
 G dd„ deƒZG d	d
„ d
eƒZG dd„ deƒZG dd„ deƒZdS )z;Perform aggregation operations on a collection or database.é    )ÚSON)Úcommon)Úvalidate_collation_or_none)ÚConfigurationError)ÚReadPreferencec               @   sn   e Zd ZdZd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S )Ú_AggregationCommandzñThe internal abstract base class for aggregation cursors.

    Should not be called directly by application developers. Use
    :meth:`pymongo.collection.Collection.aggregate`, or
    :meth:`pymongo.database.Database.aggregate` instead.
    Nc             C   sæ   d|krt dƒ‚|| _tjd|ƒ || _d| _|rPd|d ksJd|d krPd| _tjd	|ƒ || _tjd
| jj	d
d ƒƒ| _
| jjdi ƒ | j
d k	rª| j rª| j
| jd d
< || _|| _|| _|| _t|j	dd ƒƒ| _|j	dd ƒ| _d S )NZexplainzBThe explain option is not supported. Use Database.command instead.ÚpipelineFz$outé   z$mergeTÚoptionsÚ	batchSizeÚcursorÚ	collationZmaxAwaitTimeMSéÿÿÿÿr   )r   Ú_targetr   Zvalidate_listÚ	_pipelineÚ_performs_writeZvalidate_is_mappingÚ_optionsZ%validate_non_negative_integer_or_noneÚpopÚ_batch_sizeÚ
setdefaultÚ_cursor_classÚ_explicit_sessionÚ_user_fieldsÚ_result_processorr   Ú
_collationÚ_max_await_time_ms)ÚselfÚtargetZcursor_classr   r
   Úexplicit_sessionÚuser_fieldsZresult_processor© r    ú6/tmp/pip-build-20mum3z4/pymongo/pymongo/aggregation.pyÚ__init__    s,    z_AggregationCommand.__init__c             C   s   t ‚dS )z.The argument to pass to the aggregate command.N)ÚNotImplementedError)r   r    r    r!   Ú_aggregation_targetG   s    z'_AggregationCommand._aggregation_targetc             C   s   t ‚dS )z4The namespace in which the aggregate command is run.N)r#   )r   r    r    r!   Ú_cursor_namespaceL   s    z%_AggregationCommand._cursor_namespacec             C   s   t ‚dS )z5The Collection used for the aggregate command cursor.N)r#   )r   Z
cursor_docr    r    r!   Ú_cursor_collectionQ   s    z&_AggregationCommand._cursor_collectionc             C   s   t ‚dS )z:The database against which the aggregation command is run.N)r#   )r   r    r    r!   Ú	_databaseV   s    z_AggregationCommand._databasec             C   s   dS )z=Check whether the server version in-use supports aggregation.Nr    )Ú	sock_infor    r    r!   Ú_check_compat[   s    z!_AggregationCommand._check_compatc             C   s   | j r| j |||||ƒ d S )N)r   )r   ÚresultÚsessionÚserverr(   Úslave_okr    r    r!   Ú_process_result`   s    z#_AggregationCommand._process_resultc             C   s   | j rtjS | jj|ƒS )N)r   r   ZPRIMARYr   Z_read_preference_for)r   r+   r    r    r!   Úget_read_preferencee   s    z'_AggregationCommand.get_read_preferencec       
      C   s  | j |ƒ td| jfd| jfgƒ}|j| jƒ d|kr\|jdkrH| j sR|jdkr\| jj	}nd }d|kr|| jr|| jj
|ƒ}nd }|j| jj||| j|ƒ| jjd||| j|| jj| jd�}| j|||||ƒ d	|krÜ|d	 }	nd
|jdg ƒ| jdœ}	| j| j|	ƒ|	|j| j�pd
| j|| jd�S )NZ	aggregater   ZreadConcerné   é   ZwriteConcernT)Zparse_write_concern_errorÚread_concernÚwrite_concernr   r+   Úclientr   r   r   r*   )ÚidZ
firstBatchÚns)Z
batch_sizeZmax_await_time_msr+   r   )r)   r   r$   r   Úupdater   Úmax_wire_versionr   r   r2   Z_write_concern_forÚcommandr'   Únamer/   Zcodec_optionsr   r4   r   r.   Úgetr%   r   r&   Úaddressr   r   r   )
r   r+   r,   r(   r-   Úcmdr2   r3   r*   r   r    r    r!   Ú
get_cursorj   sJ    









z_AggregationCommand.get_cursor)NN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r"   Úpropertyr$   r%   r&   r'   Ústaticmethodr)   r.   r/   r>   r    r    r    r!   r      s   
&r   c                   sH   e Zd Z‡ fdd„Zedd„ ƒZedd„ ƒZdd„ Zed	d
„ ƒZ‡  Z	S )Ú_CollectionAggregationCommandc                s<   |j ddƒ}tt| ƒj||Ž || _| js8| jj dd ƒ d S )NÚ
use_cursorTr   )r   ÚsuperrE   r"   Ú_use_cursorr   )r   ÚargsÚkwargsrF   )Ú	__class__r    r!   r"   ¬   s
    z&_CollectionAggregationCommand.__init__c             C   s   | j jS )N)r   r:   )r   r    r    r!   r$   ¶   s    z1_CollectionAggregationCommand._aggregation_targetc             C   s   | j jS )N)r   Z	full_name)r   r    r    r!   r%   º   s    z/_CollectionAggregationCommand._cursor_namespacec             C   s   | j S )z5The Collection used for the aggregate command cursor.)r   )r   r   r    r    r!   r&   ¾   s    z0_CollectionAggregationCommand._cursor_collectionc             C   s   | j jS )N)r   Zdatabase)r   r    r    r!   r'   Â   s    z'_CollectionAggregationCommand._database)
r?   r@   rA   r"   rC   r$   r%   r&   r'   Ú__classcell__r    r    )rK   r!   rE   «   s
   
rE   c                   s   e Zd Z‡ fdd„Z‡  ZS )Ú _CollectionRawAggregationCommandc                s2   t t| ƒj||Ž | jr.| j r.d| jd d< d S )Nr   r   r   )rG   rM   r"   rH   r   r   )r   rI   rJ   )rK   r    r!   r"   È   s    z)_CollectionRawAggregationCommand.__init__)r?   r@   rA   r"   rL   r    r    )rK   r!   rM   Ç   s   rM   c               @   sD   e Zd Zedd„ ƒZedd„ ƒZedd„ ƒZdd„ Zed	d
„ ƒZ	dS )Ú_DatabaseAggregationCommandc             C   s   dS )Nr	   r    )r   r    r    r!   r$   Ñ   s    z/_DatabaseAggregationCommand._aggregation_targetc             C   s   d| j jf S )Nz%s.$cmd.aggregate)r   r:   )r   r    r    r!   r%   Õ   s    z-_DatabaseAggregationCommand._cursor_namespacec             C   s   | j S )N)r   )r   r    r    r!   r'   Ù   s    z%_DatabaseAggregationCommand._databasec             C   s$   |j d| jƒjddƒ\}}| j| S )z5The Collection used for the aggregate command cursor.r6   Ú.r	   )r;   r%   Úsplitr'   )r   r   Ú_Zcollnamer    r    r!   r&   Ý   s    z._DatabaseAggregationCommand._cursor_collectionc             C   s   | j dksd}t|ƒ‚d S )Né   z7Database.aggregate() is only supported on MongoDB 3.6+.)r8   r   )r(   Úerr_msgr    r    r!   r)   å   s    
z)_DatabaseAggregationCommand._check_compatN)
r?   r@   rA   rC   r$   r%   r'   r&   rD   r)   r    r    r    r!   rN   Ð   s
   rN   N)rB   Zbson.sonr   Zpymongor   Zpymongo.collationr   Zpymongo.errorsr   Zpymongo.read_preferencesr   Úobjectr   rE   rM   rN   r    r    r    r!   Ú<module>   s    	