o
    Åîï]§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r%d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   sF   | ƒ }|sdS 	 z|  ¡ }W n tjy   Y dS w |\}}||Ž  q)NFT)Ú
get_nowaitÚQueueÚEmpty)Ú	queue_refÚqÚeventÚfnÚargs© r    úB/var/www/html/env/lib/python3.10/site-packages/pymongo/topology.pyÚprocess_events_queue1   s   úùr"   c                   @   s^  e Zd ZdZdd„ Zdd„ Z		dRdd„Zd	d
„ Z		dRdd„Z	dSd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dTd*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 )UÚTopologyz*Monitor a topology of one or more servers.c                    sÄ  |j | _ |jj| _| jd u}|o| jj| _|o| jj| _d | _d | _	| js(| jr/t
j
dd�| _| jr>| j | jj| j ff¡ || _t| ¡ | ¡ |jd d |ƒ}|| _| jrottji d d d | jƒ}| j | jj|| j| j ff¡ |jD ]}| jr„| j | jj|| j ff¡ qrt| ¡ ƒ| _d| _t ¡ | _| j | j¡| _ i | _!d | _"d | _#t$ƒ | _%| js¯| jrÎ‡ fdd„}t&j't(j)d|dd�}t* +| j|j,¡‰ || _	| -¡  d | _.| jj/d uràt0| | jƒ| _.d S d S )	Néd   )ÚmaxsizeFc                      s   t ˆ ƒS ©N)r"   r    ©Úweakr    r!   Útargetv   s   z!Topology.__init__.<locals>.targetg      à?Úpymongo_events_thread)ÚintervalÚmin_intervalr)   Úname)1Ú_topology_idÚ_pool_optionsÚevent_listenersÚ
_listenersÚenabled_for_serverÚ_publish_serverÚenabled_for_topologyÚ_publish_tpÚ_eventsÚ_Topology__events_executorr   ÚputÚpublish_topology_openedÚ	_settingsr	   Úget_topology_typeÚget_server_descriptionsÚreplica_set_nameÚ_descriptionr   ÚUnknownÚ$publish_topology_description_changedÚseedsÚpublish_server_openedÚlistÚserver_descriptionsÚ_seed_addressesÚ_openedÚ	threadingÚLockÚ_lockÚcondition_classÚ
_conditionÚ_serversÚ_pidÚ_max_cluster_timer   Ú_session_poolr   ÚPeriodicExecutorr   ÚEVENTS_QUEUE_FREQUENCYÚweakrefÚrefÚcloseÚopenÚ_srv_monitorÚfqdnr   )ÚselfÚtopology_settingsÚpubÚtopology_descriptionÚ
initial_tdÚseedr)   Úexecutorr    r'   r!   Ú__init__D   sx   

ÿú
ÿþ
ÿ€
ü	ÿzTopology.__init__c                 C   s’   | j du rt ¡ | _ n$t ¡ | j kr/t d¡ | j� | j ¡  W d  ƒ n1 s*w   Y  | j� |  ¡  W d  ƒ dS 1 sBw   Y  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)	rM   ÚosÚgetpidÚwarningsÚwarnrI   rO   ÚresetÚ_ensure_opened©rX   r    r    r!   rU   Š   s   
ÿý
"ÿzTopology.openNc                    s`   |du r	ˆ j j}n|}ˆ j� ˆ  |||¡}‡ fdd„|D ƒW  d  ƒ S 1 s)w   Y  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¡‘qS r    )Úget_server_by_addressÚaddress©Ú.0Úsdrf   r    r!   Ú
<listcomp>Ã   s    ÿz+Topology.select_servers.<locals>.<listcomp>)r:   Úserver_selection_timeoutrI   Ú_select_servers_loop)rX   Úselectorrm   rh   Úserver_timeoutrD   r    rf   r!   Úselect_servers§   s   
ÿ
ÿ$üzTopology.select_serversc                 C   sœ   t ƒ }|| }| jj||| jjd�}|sG|dks||kr#t|  |¡ƒ‚|  ¡  |  ¡  | j	 
tj¡ | j ¡  t ƒ }| jj||| jjd�}|r| j ¡  |S )z7select_servers() guts. Hold the lock when calling this.)Úcustom_selectorr   )Ú_timer>   Úapply_selectorr:   Úserver_selectorr   Ú_error_messagere   Ú_request_check_allrK   Úwaitr   ÚMIN_HEARTBEAT_INTERVALÚcheck_compatible)rX   ro   Útimeoutrh   ÚnowÚend_timerD   r    r    r!   rn   Æ   s,   
ÿÿ
þð
zTopology._select_servers_loopc                 C   s   t  |  |||¡¡S )zALike select_servers, but choose a random server if several match.)ÚrandomÚchoicerq   )rX   ro   rm   rh   r    r    r!   Úselect_serverä   s   
þzTopology.select_serverc                 C   s   |   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.
        )r€   r   )rX   rh   rm   r    r    r!   Úselect_server_by_addressí   s   þz!Topology.select_server_by_addressc                 C   s´   | j }| jr|j|j }| j | jj|||j| jff¡ t	| j |ƒ| _ |  
¡  |  |j¡ | jr?| j | jj|| j | jff¡ | jrS|jtjkrS| j jtvrS| j ¡  | j ¡  dS )ziProcess a new ServerDescription on an opened topology.

        Hold the lock when calling this.
        N)r>   r3   Ú_server_descriptionsrh   r6   r8   r1   Ú"publish_server_description_changedr.   r   Ú_update_serversÚ_receive_cluster_time_no_lockÚcluster_timer5   r@   rV   Útopology_typer   r?   r
   rT   rK   Ú
notify_all)rX   Úserver_descriptionÚtd_oldÚold_server_descriptionr    r    r!   Ú_process_change  s6   ÿÿþÿþÿ
zTopology._process_changec                 C   sj   | j �( | jr| j |j¡r#|  |¡ W d  ƒ dS W d  ƒ dS W d  ƒ dS 1 s.w   Y  dS )zAProcess a new ServerDescription after an ismaster call completes.N)rI   rF   r>   Ú
has_serverrh   rŒ   )rX   r‰   r    r    r!   Ú	on_change(  s   	ÿõ	÷	"÷zTopology.on_changec                 C   sH   | j }t| j |ƒ| _ |  ¡  | jr"| j | jj|| j | jff¡ dS dS )z_Process a new seedlist on an opened topology.
        Hold the lock when calling this.
        N)	r>   r   r„   r5   r6   r8   r1   r@   r.   )rX   ÚseedlistrŠ   r    r    r!   Ú_process_srv_update8  s   ÿ
þÿzTopology._process_srv_updatec                 C   sL   | j � | jr|  |¡ W d  ƒ dS W d  ƒ dS 1 sw   Y  dS )z?Process a new list of nodes obtained from scanning SRV records.N)rI   rF   r�   )rX   r�   r    r    r!   Úon_srv_updateG  s   þ"ÿzTopology.on_srv_updatec                 C   s   | 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.
        )rL   Úget©rX   rh   r    r    r!   rg   N  s   zTopology.get_server_by_addressc                 C   s
   || j v S r&   )rL   r“   r    r    r!   r�   X  s   
zTopology.has_serverc                 C   s`   | j �# | jj}|tjkr	 W d  ƒ dS t|  ¡ ƒd jW  d  ƒ S 1 s)w   Y  dS )z!Return primary's address or None.Nr   )rI   r>   r‡   r   ÚReplicaSetWithPrimaryr   Ú_new_selectionrh   )rX   r‡   r    r    r!   Úget_primary[  s   
ý$ûzTopology.get_primaryc                 C   sp   | j �+ | jj}|tjtjfvrtƒ W  d  ƒ S tdd„ ||  ¡ ƒD ƒƒW  d  ƒ S 1 s1w   Y  dS )z+Return set of replica set member addresses.Nc                 S   s   g | ]}|j ‘qS r    )rh   ri   r    r    r!   rl   n  s    z5Topology._get_replica_set_members.<locals>.<listcomp>)rI   r>   r‡   r   r”   ÚReplicaSetNoPrimaryÚsetr•   )rX   ro   r‡   r    r    r!   Ú_get_replica_set_memberse  s   ÿü$úz!Topology._get_replica_set_membersc                 C   ó
   |   t¡S )z"Return set of secondary addresses.)r™   r   rf   r    r    r!   Úget_secondariesp  ó   
zTopology.get_secondariesc                 C   rš   )z Return set of arbiter addresses.)r™   r   rf   r    r    r!   Úget_arbiterst  rœ   zTopology.get_arbitersc                 C   ó   | j S )z1Return a document, the highest seen $clusterTime.©rN   rf   r    r    r!   Úmax_cluster_timex  ó   zTopology.max_cluster_timec                 C   s.   |r| j r|d | j d kr|| _ d S d S d S )NÚclusterTimerŸ   ©rX   r†   r    r    r!   r…   |  s   ÿ
ûz&Topology._receive_cluster_time_no_lockc                 C   s6   | j � |  |¡ W d   ƒ d S 1 sw   Y  d S r&   )rI   r…   r£   r    r    r!   Úreceive_cluster_timeŠ  s   "ÿzTopology.receive_cluster_timeé   c                 C   s@   | j � |  ¡  | j |¡ W d  ƒ dS 1 sw   Y  dS )z=Wake all monitors, wait for at least one to check its server.N)rI   rw   rK   rx   )rX   Ú	wait_timer    r    r!   Úrequest_check_allŽ  s   "þzTopology.request_check_allc                 C   sV   | j � | j |¡}|r|j ¡  W d   ƒ d S W d   ƒ d S 1 s$w   Y  d S r&   )rI   rL   r’   Úpoolrd   ©rX   rh   Úserverr    r    r!   Ú
reset_pool”  s   ý"þzTopology.reset_poolc                 C   s:   | j � | j|dd� W d  ƒ dS 1 sw   Y  dS )zgClear our pool for a server and mark it Unknown.

        Do *not* request an immediate check.
        T©r«   N)rI   Ú_reset_serverr“   r    r    r!   Úreset_serverš  s   "ÿzTopology.reset_serverc                 C   óD   | j � | j|dd� |  |¡ W d  ƒ dS 1 sw   Y  dS )z@Clear our pool for a server, mark it Unknown, and check it soon.Tr¬   N©rI   r­   Ú_request_checkr“   r    r    r!   Úreset_server_and_request_check¢  ó   "þz'Topology.reset_server_and_request_checkc                 C   r¯   )z)Mark a server Unknown, and check it soon.Fr¬   Nr°   r“   r    r    r!   Ú%mark_server_unknown_and_request_check¨  r³   z.Topology.mark_server_unknown_and_request_checkc                 C   sj   g }| j � | j ¡ D ]}| ||jjf¡ qW d   ƒ n1 s!w   Y  |D ]
\}}|j |¡ q(d S r&   )rI   rL   ÚvaluesÚappendÚ_poolÚpool_idÚremove_stale_sockets)rX   Úserversrª   r¸   r    r    r!   Úupdate_pool®  s   ÿÿÿzTopology.update_poolc                 C   sÊ   | j �< | j ¡ D ]}| ¡  q	| j ¡ | _| j ¡  ¡ D ]\}}|| jv r,|| j| _q| j	r5| j	 ¡  d| _
W d  ƒ n1 sBw   Y  | jrV| j | jj| jff¡ | js\| jrc| j ¡  dS dS )z?Clear pools and terminate monitors. Topology reopens on demand.FN)rI   rL   rµ   rT   r>   rd   rD   ÚitemsÚdescriptionrV   rF   r5   r6   r8   r1   Úpublish_topology_closedr.   r3   r7   )rX   rª   rh   rk   r    r    r!   rT   ¸  s&   

€
òÿÿzTopology.closec                 C   rž   r&   )r>   rf   r    r    r!   r½   Ñ  r¡   zTopology.descriptionc                 C   s4   | j � | j ¡ W  d  ƒ S 1 sw   Y  dS )z"Pop all session ids from the pool.N)rI   rO   Úpop_allrf   r    r    r!   Úpop_all_sessionsÕ  s   $ÿzTopology.pop_all_sessionsc                 C   s¢   | j �D | jj}|du r.| jjtjkr!| jjs |  t| j	j
d¡ n| jjs.|  t| j	j
d¡ | jj}|du r:tdƒ‚| j |¡W  d  ƒ S 1 sJw   Y  dS )z>Start or resume a server session, or raise ConfigurationError.Nz5Sessions are not supported by this MongoDB deployment)rI   r>   Úlogical_session_timeout_minutesr‡   r   ÚSingleÚhas_known_serversrn   r   r:   rm   Úreadable_serversr   r   rO   Úget_server_session)rX   Úsession_timeoutr    r    r!   rÅ   Ú  s0   ý€ýÿ
$ëzTopology.get_server_sessionc                 C   sn   |r/| j �  | jj}|d ur| j ||¡ W d   ƒ d S W d   ƒ d S 1 s(w   Y  d S | j |¡ d S r&   )rI   r>   rÁ   rO   Úreturn_server_sessionÚreturn_server_session_no_lock)rX   Úserver_sessionÚlockrÆ   r    r    r!   rÇ   ó  s   ÿÿü"ýzTopology.return_server_sessionc                 C   s   t  | j¡S )zmA Selection object, initially including all known servers.

        Hold the lock when calling this.
        )r   Úfrom_topology_descriptionr>   rf   r    r    r!   r•   ÿ  s   zTopology._new_selectionc                 C   sb   | j s#d| _ |  ¡  | js| jr| j ¡  | jr#| jjt	v r#| j ¡  t
| jƒD ]}| ¡  q(dS )z[Start monitors, or restart after a fork.

        Hold the lock when calling this.
        TN)rF   r„   r5   r3   r7   rU   rV   r½   r‡   r
   r   rL   ©rX   rª   r    r    r!   re     s   
ÿ

ÿzTopology._ensure_openedc                 C   s:   | j  |¡}|r|r| ¡  | j |¡| _|  ¡  dS 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)rL   r’   rd   r>   r®   r„   )rX   rh   r«   rª   r    r    r!   r­     s   úzTopology._reset_serverc                 C   s    | j  |¡}|r| ¡  dS dS )z2Wake one monitor. Hold the lock when calling this.N)rL   r’   Úrequest_checkr©   r    r    r!   r±   ,  s   ÿzTopology._request_checkc                 C   s   | j  ¡ D ]}| ¡  qdS )z3Wake all monitors. Hold the lock when calling this.N)rL   rµ   rÍ   rÌ   r    r    r!   rw   4  s   
ÿzTopology._request_check_allc              	   C   sú   | j  ¡  ¡ D ]W\}}|| jvrB| jj|| |  |¡| jd�}d}| jr)t 	| j
¡}t||  |¡|| j| j|d�}|| j|< | ¡  q| j| jj}|| j| _||jkr^| j| j |j¡ qt| j ¡ ƒD ]\}}| j  |¡sz| ¡  | j |¡ qfdS )zrSync our Servers from TopologyDescription.server_descriptions.

        Hold the lock while calling this.
        )r‰   Útopologyr¨   rY   N)r‰   r¨   ÚmonitorÚtopology_idÚ	listenersÚevents)r>   rD   r¼   rL   r:   Úmonitor_classÚ_create_pool_for_monitorr3   rR   rS   r6   r   Ú_create_pool_for_serverr.   r1   rU   r½   Úis_writabler¨   Úupdate_is_writablerC   r�   rT   Úpop)rX   rh   rk   rÏ   r(   rª   Úwas_writabler    r    r!   r„   9  sD   
üú


ÿ€€ýzTopology._update_serversc                 C   s   | j  || j j¡S r&   )r:   Ú
pool_classÚpool_optionsr“   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_hostnamer0   ÚappnameÚdriverF)Ú	handshake)
r:   rÛ   r   rÜ   rÞ   rß   r0   rà   rá   rÚ   )rX   rh   ÚoptionsÚmonitor_pool_optionsr    r    r!   rÔ   e  s   ù
	ÿz!Topology._create_pool_for_monitorc                    s  | j jtjtjfv }|rd}n| j jtjkrd}nd}| j jr1|tu r+|r'dS d| S d||f S t| j  	¡ ƒ}t| j  	¡  
¡ ƒ}|sQ|rMd|| jjf S d| S |d	 j‰ t‡ fd
d„|dd… D ƒƒ}|r�ˆ du rod| S |r}t|ƒ | j¡s}d| S tˆ ƒS d dd„ |D ƒ¡S )zeFormat an error message if server selection fails.

        Hold the lock when calling this.
        zreplica set membersÚ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 r&   ©Úerror©rj   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 r&   )rç   Ústrrè   r    r    r!   ré   ­  s   € ÿ)r>   r‡   r   r”   r—   ÚShardedÚknown_serversr   rC   rD   rµ   r:   r=   rç   Úallr˜   ÚintersectionrE   rì   Újoin)rX   ro   Úis_replica_setÚserver_pluralÚ	addressesrº   Úsamer    ræ   r!   rv   w  sJ   þÿ
ÿþÿzTopology._error_message)NNr&   )r¥   ),Ú__name__Ú
__module__Ú__qualname__Ú__doc__r_   rU   rq   rn   r€   r�   rŒ   rŽ   r�   r‘   rg   r�   r–   r™   r›   r�   r    r…   r¤   r§   r«   r®   r²   r´   r»   rT   Úpropertyr½   rÀ   rÅ   rÇ   r•   re   r­   r±   rw   r„   rÕ   rÔ   rv   r    r    r    r!   r#   B   s^    F
ý 
ý

ÿ$




)r#   ),rù   r`   r~   rG   rb   rR   Úbson.py3compatr   r   Úqueuer   Úpymongor   r   Úpymongo.poolr   Úpymongo.topology_descriptionr   r   r	   r
   r   Úpymongo.errorsr   r   Úpymongo.monitorr   Úpymongo.monotonicr   rs   Úpymongo.serverr   Úpymongo.server_selectorsr   r   r   r   r   r   Úpymongo.client_sessionr   r"   Úobjectr#   r    r    r    r!   Ú<module>   s,   
 