1717from __future__ import annotations
1818
1919import asyncio
20- import logging
2120import os
2221import queue
2322import random
3029from typing import TYPE_CHECKING , Any , Callable , Optional , cast
3130
3231from pymongo import _csot , common , helpers_shared , periodic_executor
33- from pymongo ._telemetry import _SdamTelemetry
32+ from pymongo ._telemetry import _SdamTelemetry , _ServerSelectionTelemetry
3433from pymongo .asynchronous .client_session import _ServerSession , _ServerSessionPool
3534from pymongo .asynchronous .monitor import MonitorBase , SrvMonitor
3635from pymongo .asynchronous .pool import Pool
5251 _async_create_condition ,
5352 _async_create_lock ,
5453)
55- from pymongo .logger import (
56- _SERVER_SELECTION_LOGGER ,
57- _debug_log ,
58- _ServerSelectionStatusMessage ,
59- )
6054from pymongo .pool_options import PoolOptions
6155from pymongo .server_description import ServerDescription
6256from pymongo .server_selectors import (
@@ -107,14 +101,14 @@ class Topology:
107101 def __init__ (self , topology_settings : TopologySettings ):
108102 self ._topology_id = topology_settings ._topology_id
109103 self ._listeners = topology_settings ._pool_options ._event_listeners
110- self ._publish_server = self ._listeners is not None and self ._listeners .enabled_for_server
111- self ._publish_tp = self ._listeners is not None and self ._listeners .enabled_for_topology
112104
113105 # Create events queue if there are publishers.
114106 self ._events : queue .Queue [Any ] | None = None
115107 self .__events_executor : Any = None
116108
117- if self ._publish_server or self ._publish_tp :
109+ publish_server = self ._listeners is not None and self ._listeners .enabled_for_server
110+ publish_tp = self ._listeners is not None and self ._listeners .enabled_for_topology
111+ if publish_server or publish_tp :
118112 self ._events = queue .Queue (maxsize = 100 )
119113
120114 self ._sdam = _SdamTelemetry (self ._topology_id , self ._listeners , self ._events )
@@ -152,7 +146,7 @@ def __init__(self, topology_settings: TopologySettings):
152146 self ._max_cluster_time : Optional [ClusterTime ] = None
153147 self ._session_pool = _ServerSessionPool ()
154148
155- if self ._publish_server or self ._publish_tp :
149+ if self ._sdam . _publish_server or self . _sdam ._publish_tp :
156150 assert self ._events is not None
157151 weak : weakref .ReferenceType [queue .Queue [Any ]]
158152
@@ -285,17 +279,10 @@ async def _select_servers_loop(
285279 now = time .monotonic ()
286280 end_time = now + timeout
287281 logged_waiting = False
288-
289- if _SERVER_SELECTION_LOGGER .isEnabledFor (logging .DEBUG ):
290- _debug_log (
291- _SERVER_SELECTION_LOGGER ,
292- message = _ServerSelectionStatusMessage .STARTED ,
293- selector = selector ,
294- operation = operation ,
295- operationId = operation_id ,
296- topologyDescription = self .description ,
297- clientId = self .description ._topology_settings ._topology_id ,
298- )
282+ ss = _ServerSelectionTelemetry (
283+ self ._topology_id , selector , operation , operation_id , self .description
284+ )
285+ ss .started ()
299286
300287 server_descriptions = self ._description .apply_selector (
301288 selector ,
@@ -309,32 +296,13 @@ async def _select_servers_loop(
309296 while not server_descriptions :
310297 # No suitable servers.
311298 if timeout == 0 or now > end_time :
312- if _SERVER_SELECTION_LOGGER .isEnabledFor (logging .DEBUG ):
313- _debug_log (
314- _SERVER_SELECTION_LOGGER ,
315- message = _ServerSelectionStatusMessage .FAILED ,
316- selector = selector ,
317- operation = operation ,
318- operationId = operation_id ,
319- topologyDescription = self .description ,
320- clientId = self .description ._topology_settings ._topology_id ,
321- failure = self ._error_message (selector ),
322- )
299+ ss .failed (self ._error_message (selector ))
323300 raise ServerSelectionTimeoutError (
324301 f"{ self ._error_message (selector )} , Timeout: { timeout } s, Topology Description: { self .description !r} "
325302 )
326303
327304 if not logged_waiting :
328- _debug_log (
329- _SERVER_SELECTION_LOGGER ,
330- message = _ServerSelectionStatusMessage .WAITING ,
331- selector = selector ,
332- operation = operation ,
333- operationId = operation_id ,
334- topologyDescription = self .description ,
335- clientId = self .description ._topology_settings ._topology_id ,
336- remainingTimeMS = int (1000 * (end_time - time .monotonic ())),
337- )
305+ ss .waiting (int (1000 * (end_time - time .monotonic ())))
338306 logged_waiting = True
339307
340308 await self ._ensure_opened ()
@@ -399,18 +367,9 @@ async def select_server(
399367 )
400368 if _csot .get_timeout ():
401369 _csot .set_rtt (server .description .min_round_trip_time )
402- if _SERVER_SELECTION_LOGGER .isEnabledFor (logging .DEBUG ):
403- _debug_log (
404- _SERVER_SELECTION_LOGGER ,
405- message = _ServerSelectionStatusMessage .SUCCEEDED ,
406- selector = selector ,
407- operation = operation ,
408- operationId = operation_id ,
409- topologyDescription = self .description ,
410- clientId = self .description ._topology_settings ._topology_id ,
411- serverHost = server .description .address [0 ],
412- serverPort = server .description .address [1 ],
413- )
370+ _ServerSelectionTelemetry (
371+ self ._topology_id , selector , operation , operation_id , self .description
372+ ).succeeded (server .description .address [0 ], server .description .address [1 ])
414373 return server
415374
416375 async def select_server_by_address (
@@ -670,7 +629,7 @@ async def close(self) -> None:
670629 self ._closed = True
671630
672631 # Publish only after releasing the lock.
673- if self ._publish_tp :
632+ if self ._sdam . _publish_tp :
674633 self ._description = TopologyDescription (
675634 TOPOLOGY_TYPE .Unknown ,
676635 {},
@@ -681,7 +640,7 @@ async def close(self) -> None:
681640 )
682641 self ._sdam .topology_closed (old_td , self ._description )
683642
684- if self ._publish_server or self ._publish_tp :
643+ if self ._sdam . _publish_server or self . _sdam ._publish_tp :
685644 # Make sure the events executor thread is fully closed before publishing the remaining events
686645 self .__events_executor .close ()
687646 await self .__events_executor .join (1 )
@@ -722,7 +681,7 @@ async def _ensure_opened(self) -> None:
722681 await self ._update_servers ()
723682
724683 # Start or restart the events publishing thread.
725- if self ._publish_tp or self ._publish_server :
684+ if self ._sdam . _publish_tp or self . _sdam ._publish_server :
726685 self .__events_executor .open ()
727686
728687 # Start the SRV polling thread.
@@ -861,7 +820,7 @@ async def _update_servers(self) -> None:
861820 )
862821
863822 weak = None
864- if self ._publish_server and self ._events is not None :
823+ if self ._sdam . _publish_server and self ._events is not None :
865824 weak = weakref .ref (self ._events )
866825 server = Server (
867826 server_description = sd ,
0 commit comments