Skip to content

Commit 577b594

Browse files
committed
PYTHON-5846 Add _emit_log helpers to _HeartbeatTelemetry and _SdamTelemetry
1 parent a30d380 commit 577b594

1 file changed

Lines changed: 57 additions & 75 deletions

File tree

pymongo/_telemetry.py

Lines changed: 57 additions & 75 deletions
Original file line numberDiff line numberDiff line change
@@ -359,6 +359,18 @@ def __init__(
359359
self._publish = publish
360360
self._awaited = awaited
361361

362+
def _emit_log(self, message: _SDAMStatusMessage, **extra: Any) -> None:
363+
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
364+
_debug_log(
365+
_SDAM_LOGGER,
366+
message=message,
367+
topologyId=self._topology_id,
368+
serverHost=self._address[0],
369+
serverPort=self._address[1],
370+
awaited=self._awaited,
371+
**extra,
372+
)
373+
362374
def started(self) -> None:
363375
"""Publish the APM heartbeat-started event (before connection checkout)."""
364376
self._start = time.monotonic()
@@ -368,17 +380,11 @@ def started(self) -> None:
368380

369381
def emit_started_log(self, conn_id: int, server_conn_id: Optional[int]) -> None:
370382
"""Emit the log entry for heartbeat started (after connection checkout)."""
371-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
372-
_debug_log(
373-
_SDAM_LOGGER,
374-
message=_SDAMStatusMessage.HEARTBEAT_START,
375-
topologyId=self._topology_id,
376-
driverConnectionId=conn_id,
377-
serverConnectionId=server_conn_id,
378-
serverHost=self._address[0],
379-
serverPort=self._address[1],
380-
awaited=self._awaited,
381-
)
383+
self._emit_log(
384+
_SDAMStatusMessage.HEARTBEAT_START,
385+
driverConnectionId=conn_id,
386+
serverConnectionId=server_conn_id,
387+
)
382388

383389
def succeeded(
384390
self,
@@ -393,19 +399,13 @@ def succeeded(
393399
self._listeners.publish_server_heartbeat_succeeded(
394400
self._address, round_trip_time, response, response.awaitable
395401
)
396-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
397-
_debug_log(
398-
_SDAM_LOGGER,
399-
message=_SDAMStatusMessage.HEARTBEAT_SUCCESS,
400-
topologyId=self._topology_id,
401-
driverConnectionId=conn_id,
402-
serverConnectionId=server_conn_id,
403-
serverHost=self._address[0],
404-
serverPort=self._address[1],
405-
awaited=self._awaited,
406-
durationMS=round_trip_time * 1000,
407-
reply=response.document,
408-
)
402+
self._emit_log(
403+
_SDAMStatusMessage.HEARTBEAT_SUCCESS,
404+
driverConnectionId=conn_id,
405+
serverConnectionId=server_conn_id,
406+
durationMS=round_trip_time * 1000,
407+
reply=response.document,
408+
)
409409

410410
def failed(self, error: Exception, conn_id: Optional[int]) -> None:
411411
"""Emit the FAILED log entry and APM event."""
@@ -415,18 +415,12 @@ def failed(self, error: Exception, conn_id: Optional[int]) -> None:
415415
self._listeners.publish_server_heartbeat_failed(
416416
self._address, duration, error, self._awaited
417417
)
418-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
419-
_debug_log(
420-
_SDAM_LOGGER,
421-
message=_SDAMStatusMessage.HEARTBEAT_FAIL,
422-
topologyId=self._topology_id,
423-
serverHost=self._address[0],
424-
serverPort=self._address[1],
425-
awaited=self._awaited,
426-
durationMS=duration * 1000,
427-
failure=error,
428-
driverConnectionId=conn_id,
429-
)
418+
self._emit_log(
419+
_SDAMStatusMessage.HEARTBEAT_FAIL,
420+
durationMS=duration * 1000,
421+
failure=error,
422+
driverConnectionId=conn_id,
423+
)
430424

431425

432426
class _SdamTelemetry:
@@ -453,13 +447,17 @@ def _enqueue(self, fn: Any, args: tuple[Any, ...]) -> None:
453447
if self._events is not None:
454448
self._events.put((fn, args))
455449

456-
def topology_opened(self) -> None:
450+
def _emit_log(self, message: _SDAMStatusMessage, **extra: Any) -> None:
457451
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
458452
_debug_log(
459453
_SDAM_LOGGER,
460-
message=_SDAMStatusMessage.START_TOPOLOGY,
454+
message=message,
461455
topologyId=self._topology_id,
456+
**extra,
462457
)
458+
459+
def topology_opened(self) -> None:
460+
self._emit_log(_SDAMStatusMessage.START_TOPOLOGY)
463461
if self._publish_tp:
464462
assert self._listeners is not None
465463
self._enqueue(self._listeners.publish_topology_opened, (self._topology_id,))
@@ -471,14 +469,11 @@ def topology_description_changed(self, old_td: Any, new_td: Any) -> None:
471469
self._listeners.publish_topology_description_changed,
472470
(old_td, new_td, self._topology_id),
473471
)
474-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
475-
_debug_log(
476-
_SDAM_LOGGER,
477-
message=_SDAMStatusMessage.TOPOLOGY_CHANGE,
478-
topologyId=self._topology_id,
479-
previousDescription=repr(old_td),
480-
newDescription=repr(new_td),
481-
)
472+
self._emit_log(
473+
_SDAMStatusMessage.TOPOLOGY_CHANGE,
474+
previousDescription=repr(old_td),
475+
newDescription=repr(new_td),
476+
)
482477

483478
def topology_closed(self, old_td: Any, new_td: Any) -> None:
484479
"""Emit APM and log events for topology description change + topology closed."""
@@ -489,32 +484,22 @@ def topology_closed(self, old_td: Any, new_td: Any) -> None:
489484
(old_td, new_td, self._topology_id),
490485
)
491486
self._enqueue(self._listeners.publish_topology_closed, (self._topology_id,))
492-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
493-
_debug_log(
494-
_SDAM_LOGGER,
495-
message=_SDAMStatusMessage.TOPOLOGY_CHANGE,
496-
topologyId=self._topology_id,
497-
previousDescription=repr(old_td),
498-
newDescription=repr(new_td),
499-
)
500-
_debug_log(
501-
_SDAM_LOGGER,
502-
message=_SDAMStatusMessage.STOP_TOPOLOGY,
503-
topologyId=self._topology_id,
504-
)
487+
self._emit_log(
488+
_SDAMStatusMessage.TOPOLOGY_CHANGE,
489+
previousDescription=repr(old_td),
490+
newDescription=repr(new_td),
491+
)
492+
self._emit_log(_SDAMStatusMessage.STOP_TOPOLOGY)
505493

506494
def server_opened(self, address: _Address) -> None:
507495
if self._publish_server:
508496
assert self._listeners is not None
509497
self._enqueue(self._listeners.publish_server_opened, (address, self._topology_id))
510-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
511-
_debug_log(
512-
_SDAM_LOGGER,
513-
message=_SDAMStatusMessage.START_SERVER,
514-
topologyId=self._topology_id,
515-
serverHost=address[0],
516-
serverPort=address[1],
517-
)
498+
self._emit_log(
499+
_SDAMStatusMessage.START_SERVER,
500+
serverHost=address[0],
501+
serverPort=address[1],
502+
)
518503

519504
def server_description_changed(self, sd_old: Any, sd_new: Any, address: _Address) -> None:
520505
if self._publish_server:
@@ -528,11 +513,8 @@ def server_closed(self, address: _Address) -> None:
528513
if self._publish_server:
529514
assert self._listeners is not None
530515
self._enqueue(self._listeners.publish_server_closed, (address, self._topology_id))
531-
if _SDAM_LOGGER.isEnabledFor(logging.DEBUG):
532-
_debug_log(
533-
_SDAM_LOGGER,
534-
message=_SDAMStatusMessage.STOP_SERVER,
535-
topologyId=self._topology_id,
536-
serverHost=address[0],
537-
serverPort=address[1],
538-
)
516+
self._emit_log(
517+
_SDAMStatusMessage.STOP_SERVER,
518+
serverHost=address[0],
519+
serverPort=address[1],
520+
)

0 commit comments

Comments
 (0)