3
]ð]§l  ã               @   s  d Z ddlZddlZddlZddlZddlZddlmZmZ erJddl	Z
nddl
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mZ ddlmZ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$m%Z%m&Z& ddl'm(Z( dd„ Z)G dd„ de*ƒZ+dS )z<Internal class to monitor a topology of one or more servers.é    N)Ú
itervaluesÚPY3)Úcommon)Úperiodic_executor)ÚPoolOptions)Úupdated_topology_descriptionÚ)_updated_topology_description_srv_pollingÚTopologyDescriptionÚSRV_POLLING_TOPOLOGIESÚTOPOLOGY_TYPE)ÚServerSelectionTimeoutErrorÚConfigurationError)Ú
SrvMonitor)Útime)ÚServer)Úany_server_selectorÚarbiter_server_selectorÚsecondary_server_selectorÚreadable_server_selectorÚwritable_server_selectorÚ	Selection)Ú_ServerSessionPoolc             C   sN   | ƒ }|sdS x:y|j ƒ }W n tjk
r4   P Y qX |\}}||Ž  qW dS )NFT)Ú
get_nowaitÚQueueÚEmpty)Z	queue_refÚqÚeventÚfnÚargs© r   ú3/tmp/pip-build-20mum3z4/pymongo/pymongo/topology.pyÚprocess_events_queue1   s    r!   c               @   sT  e Zd ZdZdd„ Zdd„ ZdRdd„Zd	d
„ ZdSdd„ZdTd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!d"„ Zd#d$„ Zd%d&„ Zd'd(„ ZdUd*d+„Zd,d-„ Zd.d/„ Zd0d1„ Zd2d3„ Zd4d5„ Zd6d7„ Zed8d9„ ƒZd:d;„ Z d<d=„ Z!d>d?„ Z"d@dA„ Z#dBdC„ Z$dDdE„ Z%dFdG„ Z&dHdI„ Z'dJdK„ Z(dLdM„ Z)dNdO„ Z*dPdQ„ Z+dS )VÚTopologyz*Monitor a topology of one or more servers.c                sÊ  |j | _ |jj| _| jd k	}|o&| jj| _|o4| jj| _d | _d | _	| jsP| jr^t
j
dd�| _| jr|| jj| jj| j ffƒ || _t|jƒ |jƒ |jd d |ƒ}|| _| jrÞttji d d d | jƒ}| jj| jj|| j| j ffƒ x.|jD ]$}| jræ| jj| jj|| j ffƒ qæW t|jƒ ƒ| _d| _tjƒ | _| jj| jƒ| _ i | _!d | _"d | _#t$ƒ | _%| j�sf| j�r¤‡ fdd„}t&j't(j)d|dd�}t*j+| j|j,ƒ‰ || _	|j-ƒ  d | _.| jj/d k	�rÆt0| | jƒ| _.d S )	Néd   )ÚmaxsizeFc                  s   t ˆ ƒS )N)r!   r   )Úweakr   r    Útargetv   s    z!Topology.__init__.<locals>.targetg      à?Zpymongo_events_thread)ÚintervalZmin_intervalr&   Úname)1Ú_topology_idZ_pool_optionsÚevent_listenersÚ
_listenersZenabled_for_serverÚ_publish_serverZenabled_for_topologyÚ_publish_tpÚ_eventsÚ_Topology__events_executorr   ÚputZpublish_topology_openedÚ	_settingsr	   Zget_topology_typeZget_server_descriptionsÚreplica_set_nameÚ_descriptionr   ÚUnknownÚ$publish_topology_description_changedZseedsZpublish_server_openedÚlistÚserver_descriptionsÚ_seed_addressesÚ_openedÚ	threadingÚLockÚ_lockZcondition_classÚ
_conditionÚ_serversÚ_pidÚ_max_cluster_timer   Ú_session_poolr   ZPeriodicExecutorr   ZEVENTS_QUEUE_FREQUENCYÚweakrefÚrefÚcloseÚopenÚ_srv_monitorZfqdnr   )ÚselfÚtopology_settingsZpubZtopology_descriptionZ
initial_tdÚseedr&   Úexecutorr   )r%   r    Ú__init__D   sh    



zTopology.__init__c             C   sh   | j dkrtjƒ | _ n4tjƒ | j krJtjdƒ | j� | jjƒ  W dQ R X | j� | jƒ  W dQ R X dS )a³  Start monitoring, or restart after a fork.

        No effect if called multiple times.

        .. warning:: Topology is shared among multiple threads and is protected
          by mutual exclusion. Using Topology from a process other than the one
          that initialized it will emit a warning and may result in deadlock. To
          prevent this from happening, MongoClient must be created after any
          forking.

        Nz³MongoClient opened before fork. Create MongoClient only after forking. See PyMongo's documentation for details: http://api.mongodb.org/python/current/faq.html#is-pymongo-fork-safe)	r?   ÚosÚgetpidÚwarningsÚwarnr<   rA   ÚresetÚ_ensure_opened)rG   r   r   r    rE   Š   s    
zTopology.openNc                sH   |dkrˆ j j}n|}ˆ j�" ˆ j|||ƒ}‡ fdd„|D ƒS Q R X dS )aL  Return a list of Servers matching selector, or time out.

        :Parameters:
          - `selector`: function that takes a list of Servers and returns
            a subset of them.
          - `server_selection_timeout` (optional): maximum seconds to wait.
            If not provided, the default value common.SERVER_SELECTION_TIMEOUT
            is used.
          - `address`: optional server address to select.

        Calls self.open() if needed.

        Raises exc:`ServerSelectionTimeoutError` after
        `server_selection_timeout` if no matching servers are found.
        Nc                s   g | ]}ˆ j |jƒ‘qS r   )Úget_server_by_addressÚaddress)Ú.0Úsd)rG   r   r    ú
<listcomp>Ã   s   z+Topology.select_servers.<locals>.<listcomp>)r1   Úserver_selection_timeoutr<   Ú_select_servers_loop)rG   ÚselectorrW   rS   Zserver_timeoutr7   r   )rG   r    Úselect_servers§   s    


zTopology.select_serversc             C   sž   t ƒ }|| }| jj||| jjd�}xj|sŽ|dks:||krHt| j|ƒƒ‚| jƒ  | jƒ  | j	j
tjƒ | jjƒ  t ƒ }| jj||| jjd�}q&W | jjƒ  |S )z7select_servers() guts. Hold the lock when calling this.)Zcustom_selectorr   )Ú_timer3   Zapply_selectorr1   Zserver_selectorr   Ú_error_messagerQ   Ú_request_check_allr=   Úwaitr   ZMIN_HEARTBEAT_INTERVALZcheck_compatible)rG   rY   ÚtimeoutrS   ÚnowÚend_timer7   r   r   r    rX   Æ   s$    

zTopology._select_servers_loopc             C   s   t j| j|||ƒƒS )zALike select_servers, but choose a random server if several match.)ÚrandomÚchoicerZ   )rG   rY   rW   rS   r   r   r    Úselect_serverä   s    
zTopology.select_serverc             C   s   | j t||ƒS )a‰  Return a Server for "address", reconnecting if necessary.

        If the server's type is not known, request an immediate check of all
        servers. Time out after "server_selection_timeout" if the server
        cannot be reached.

        :Parameters:
          - `address`: A (host, port) pair.
          - `server_selection_timeout` (optional): maximum seconds to wait.
            If not provided, the default value
            common.SERVER_SELECTION_TIMEOUT is used.

        Calls self.open() if needed.

        Raises exc:`ServerSelectionTimeoutError` after
        `server_selection_timeout` if no matching servers are found.
        )rd   r   )rG   rS   rW   r   r   r    Úselect_server_by_addressí   s    z!Topology.select_server_by_addressc             C   s´   | j }| jr8|j|j }| jj| jj|||j| jffƒ t	| j |ƒ| _ | j
ƒ  | j|jƒ | jr~| jj| jj|| j | jffƒ | jr¦|jtjkr¦| j jtkr¦| jjƒ  | jjƒ  dS )ziProcess a new ServerDescription on an opened topology.

        Hold the lock when calling this.
        N)r3   r,   Z_server_descriptionsrS   r.   r0   r+   Z"publish_server_description_changedr)   r   Ú_update_serversÚ_receive_cluster_time_no_lockÚcluster_timer-   r5   rF   Útopology_typer   r4   r
   rD   r=   Ú
notify_all)rG   Úserver_descriptionÚtd_oldZold_server_descriptionr   r   r    Ú_process_change  s*    
zTopology._process_changec          	   C   s4   | j �$ | jr&| jj|jƒr&| j|ƒ W dQ R X dS )zAProcess a new ServerDescription after an ismaster call completes.N)r<   r9   r3   Ú
has_serverrS   rm   )rG   rk   r   r   r    Ú	on_change(  s    	zTopology.on_changec             C   sD   | j }t| j |ƒ| _ | jƒ  | jr@| jj| jj|| j | jffƒ dS )z_Process a new seedlist on an opened topology.
        Hold the lock when calling this.
        N)	r3   r   rf   r-   r.   r0   r+   r5   r)   )rG   Úseedlistrl   r   r   r    Ú_process_srv_update8  s    zTopology._process_srv_updatec          	   C   s&   | j � | jr| j|ƒ W dQ R X dS )z?Process a new list of nodes obtained from scanning SRV records.N)r<   r9   rq   )rG   rp   r   r   r    Úon_srv_updateG  s    zTopology.on_srv_updatec             C   s   | j j|ƒS )aJ  Get a Server or None.

        Returns the current version of the server immediately, even if it's
        Unknown or absent from the topology. Only use this in unittests.
        In driver code, use select_server_by_address, since then you're
        assured a recent view of the server's type and wire protocol version.
        )r>   Úget)rG   rS   r   r   r    rR   N  s    zTopology.get_server_by_addressc             C   s
   || j kS )N)r>   )rG   rS   r   r   r    rn   X  s    zTopology.has_serverc          	   C   s:   | j �* | jj}|tjkrdS t| jƒ ƒd jS Q R X dS )z!Return primary's address or None.Nr   )r<   r3   ri   r   ÚReplicaSetWithPrimaryr   Ú_new_selectionrS   )rG   ri   r   r   r    Úget_primary[  s
    
zTopology.get_primaryc             C   sJ   | j �: | jj}|tjtjfkr&tƒ S tdd„ || jƒ ƒD ƒƒS Q R X dS )z+Return set of replica set member addresses.c             S   s   g | ]
}|j ‘qS r   )rS   )rT   rU   r   r   r    rV   n  s    z5Topology._get_replica_set_members.<locals>.<listcomp>N)r<   r3   ri   r   rt   ÚReplicaSetNoPrimaryÚsetru   )rG   rY   ri   r   r   r    Ú_get_replica_set_memberse  s    
z!Topology._get_replica_set_membersc             C   s
   | j tƒS )z"Return set of secondary addresses.)ry   r   )rG   r   r   r    Úget_secondariesp  s    zTopology.get_secondariesc             C   s
   | j tƒS )z Return set of arbiter addresses.)ry   r   )rG   r   r   r    Úget_arbiterst  s    zTopology.get_arbitersc             C   s   | j S )z1Return a document, the highest seen $clusterTime.)r@   )rG   r   r   r    Úmax_cluster_timex  s    zTopology.max_cluster_timec             C   s(   |r$| j  s|d | j d kr$|| _ d S )NZclusterTime)r@   )rG   rh   r   r   r    rg   |  s
    z&Topology._receive_cluster_time_no_lockc          	   C   s    | j � | j|ƒ W d Q R X d S )N)r<   rg   )rG   rh   r   r   r    Úreceive_cluster_timeŠ  s    zTopology.receive_cluster_timeé   c          	   C   s*   | j � | jƒ  | jj|ƒ W dQ R X dS )z=Wake all monitors, wait for at least one to check its server.N)r<   r]   r=   r^   )rG   Z	wait_timer   r   r    Úrequest_check_allŽ  s    zTopology.request_check_allc          	   C   s0   | j �  | jj|ƒ}|r"|jjƒ  W d Q R X d S )N)r<   r>   rs   ÚpoolrP   )rG   rS   Úserverr   r   r    Ú
reset_pool”  s    zTopology.reset_poolc             C   s$   | j � | j|dd� W dQ R X dS )zgClear our pool for a server and mark it Unknown.

        Do *not* request an immediate check.
        T)r‚   N)r<   Ú_reset_server)rG   rS   r   r   r    Úreset_serverš  s    zTopology.reset_serverc             C   s.   | j � | j|dd� | j|ƒ W dQ R X dS )z@Clear our pool for a server, mark it Unknown, and check it soon.T)r‚   N)r<   rƒ   Ú_request_check)rG   rS   r   r   r    Úreset_server_and_request_check¢  s    z'Topology.reset_server_and_request_checkc             C   s.   | j � | j|dd� | j|ƒ W dQ R X dS )z)Mark a server Unknown, and check it soon.F)r‚   N)r<   rƒ   r…   )rG   rS   r   r   r    Ú%mark_server_unknown_and_request_check¨  s    z.Topology.mark_server_unknown_and_request_checkc             C   s^   g }| j �, x$| jjƒ D ]}|j||jjfƒ qW W d Q R X x|D ]\}}|jj|ƒ qBW d S )N)r<   r>   ÚvaluesÚappendÚ_poolÚpool_idZremove_stale_sockets)rG   Úserversr�   r‹   r   r   r    Úupdate_pool®  s     zTopology.update_poolc             C   sº   | j �v x| jjƒ D ]}|jƒ  qW | jjƒ | _x0| jjƒ jƒ D ]\}}|| jkr@|| j| _q@W | j	rr| j	jƒ  d| _
W dQ R X | jr | jj| jj| jffƒ | js¬| jr¶| jjƒ  dS )z?Clear pools and terminate monitors. Topology reopens on demand.FN)r<   r>   rˆ   rD   r3   rP   r7   ÚitemsÚdescriptionrF   r9   r-   r.   r0   r+   Zpublish_topology_closedr)   r,   r/   )rG   r�   rS   rU   r   r   r    rD   ¸  s    

zTopology.closec             C   s   | j S )N)r3   )rG   r   r   r    r�   Ñ  s    zTopology.descriptionc          	   C   s   | j � | jjƒ S Q R X dS )z"Pop all session ids from the pool.N)r<   rA   Úpop_all)rG   r   r   r    Úpop_all_sessionsÕ  s    zTopology.pop_all_sessionsc             C   sŠ   | j �z | jj}|dkr\| jjtjkrB| jjs\| jt| j	j
dƒ n| jjs\| jt| j	j
dƒ | jj}|dkrttdƒ‚| jj|ƒS Q R X dS )z>Start or resume a server session, or raise ConfigurationError.Nz5Sessions are not supported by this MongoDB deployment)r<   r3   Úlogical_session_timeout_minutesri   r   ÚSingleZhas_known_serversrX   r   r1   rW   Zreadable_serversr   r   rA   Úget_server_session)rG   Úsession_timeoutr   r   r    r”   Ú  s&    zTopology.get_server_sessionc          
   C   sF   |r6| j �$ | jj}|d k	r*| jj||ƒ W d Q R X n| jj|ƒ d S )N)r<   r3   r’   rA   Úreturn_server_sessionZreturn_server_session_no_lock)rG   Zserver_sessionÚlockr•   r   r   r    r–   ó  s    zTopology.return_server_sessionc             C   s   t j| jƒS )zmA Selection object, initially including all known servers.

        Hold the lock when calling this.
        )r   Zfrom_topology_descriptionr3   )rG   r   r   r    ru   ÿ  s    zTopology._new_selectionc             C   sf   | j sFd| _ | jƒ  | js | jr*| jjƒ  | jrF| jjt	krF| jjƒ  xt
| jƒD ]}|jƒ  qRW dS )z[Start monitors, or restart after a fork.

        Hold the lock when calling this.
        TN)r9   rf   r-   r,   r/   rE   rF   r�   ri   r
   r   r>   )rG   r�   r   r   r    rQ     s    

zTopology._ensure_openedc             C   s6   | j j|ƒ}|r2|r|jƒ  | jj|ƒ| _| jƒ  dS )z�Mark a server Unknown and optionally reset it's pool.

        Hold the lock when calling this. Does *not* request an immediate check.
        N)r>   rs   rP   r3   r„   rf   )rG   rS   r‚   r�   r   r   r    rƒ     s    zTopology._reset_serverc             C   s   | j j|ƒ}|r|jƒ  dS )z2Wake one monitor. Hold the lock when calling this.N)r>   rs   Úrequest_check)rG   rS   r�   r   r   r    r…   ,  s    zTopology._request_checkc             C   s    x| j jƒ D ]}|jƒ  qW dS )z3Wake all monitors. Hold the lock when calling this.N)r>   rˆ   r˜   )rG   r�   r   r   r    r]   4  s    zTopology._request_check_allc          	   C   s  xÀ| j jƒ jƒ D ]®\}}|| jkr†| jj|| | j|ƒ| jd�}d}| jrTtj	| j
ƒ}t|| j|ƒ|| j| j|d�}|| j|< |jƒ  q| j| jj}|| j| _||jkr| j| jj|jƒ qW x:t| jjƒ ƒD ](\}}| j j|ƒsÒ|jƒ  | jj|ƒ qÒW dS )zrSync our Servers from TopologyDescription.server_descriptions.

        Hold the lock while calling this.
        )rk   Ztopologyr€   rH   N)rk   r€   ÚmonitorZtopology_idZ	listenersÚevents)r3   r7   rŽ   r>   r1   Zmonitor_classÚ_create_pool_for_monitorr,   rB   rC   r.   r   Ú_create_pool_for_serverr)   r+   rE   r�   Úis_writabler€   Zupdate_is_writabler6   rn   rD   Úpop)rG   rS   rU   r™   r%   r�   Zwas_writabler   r   r    rf   9  s8    




zTopology._update_serversc             C   s   | j j|| j jƒS )N)r1   Ú
pool_classÚpool_options)rG   rS   r   r   r    rœ   b  s    z Topology._create_pool_for_serverc          	   C   s>   | j j}t|j|j|j|j|j|j|jd�}| j j	||dd�S )N)Úconnect_timeoutÚsocket_timeoutÚssl_contextÚssl_match_hostnamer*   ÚappnameÚdriverF)Z	handshake)
r1   r    r   r¡   r£   r¤   r*   r¥   r¦   rŸ   )rG   rS   ÚoptionsZmonitor_pool_optionsr   r   r    r›   e  s    

z!Topology._create_pool_for_monitorc                s$  | j jtjtjfk}|rd}n| j jtjkr2d}nd}| j jrf|tkrX|rNdS d| S nd||f S nºt| j j	ƒ ƒ}t| j j	ƒ j
ƒ ƒ}|s¦|ržd|| jjf S d| S |d	 j‰ t‡ fd
d„|dd… D ƒƒ}|�rˆ dkräd| S |oøt|ƒj| jƒ �rd| S tˆ ƒS djdd„ |D ƒƒS dS )zeFormat an error message if server selection fails.

        Hold the lock when calling this.
        zreplica set membersZmongosesrŒ   zNo primary available for writeszNo %s available for writeszNo %s match selector "%s"z)No %s available for replica set name "%s"zNo %s availabler   c             3   s   | ]}|j ˆ kV  qd S )N)Úerror)rT   r�   )r¨   r   r    ú	<genexpr>�  s    z*Topology._error_message.<locals>.<genexpr>é   NzNo %s found yetz\Could not reach any servers in %s. Replica set is configured with internal hostnames or IPs?ú,c             s   s   | ]}|j rt|j ƒV  qd S )N)r¨   Ústr)rT   r�   r   r   r    r©   ­  s    )r3   ri   r   rt   rw   ZShardedZknown_serversr   r6   r7   rˆ   r1   r2   r¨   Úallrx   Úintersectionr8   r¬   Újoin)rG   rY   Zis_replica_setZserver_pluralÚ	addressesrŒ   Zsamer   )r¨   r    r\   w  s@    


zTopology._error_message)NN)NN)N)r~   ),Ú__name__Ú
__module__Ú__qualname__Ú__doc__rK   rE   rZ   rX   rd   re   rm   ro   rq   rr   rR   rn   rv   ry   rz   r{   r|   rg   r}   r   r‚   r„   r†   r‡   r�   rD   Úpropertyr�   r‘   r”   r–   ru   rQ   rƒ   r…   r]   rf   rœ   r›   r\   r   r   r   r    r"   B   sT   F 
  

$



)r"   ),r´   rL   rb   r:   rN   rB   Zbson.py3compatr   r   Úqueuer   Zpymongor   r   Zpymongo.poolr   Zpymongo.topology_descriptionr   r   r	   r
   r   Zpymongo.errorsr   r   Zpymongo.monitorr   Zpymongo.monotonicr   r[   Zpymongo.serverr   Zpymongo.server_selectorsr   r   r   r   r   r   Zpymongo.client_sessionr   r!   Úobjectr"   r   r   r   r    Ú<module>   s*   
 