|
18 | 18 |
|
19 | 19 | import datetime |
20 | 20 | import logging |
| 21 | +import queue |
21 | 22 | from collections.abc import MutableMapping |
22 | 23 | from typing import TYPE_CHECKING, Any, Optional |
23 | 24 |
|
24 | | -from pymongo.logger import _COMMAND_LOGGER, _CommandStatusMessage, _debug_log |
| 25 | +from pymongo.logger import ( |
| 26 | + _COMMAND_LOGGER, |
| 27 | + _CONNECTION_LOGGER, |
| 28 | + _SDAM_LOGGER, |
| 29 | + _CommandStatusMessage, |
| 30 | + _ConnectionStatusMessage, |
| 31 | + _debug_log, |
| 32 | + _SDAMStatusMessage, |
| 33 | + _verbose_connection_error_reason, |
| 34 | +) |
25 | 35 | from pymongo.pool_shared import _ConnectionTelemetryInfo |
26 | 36 |
|
27 | 37 | if TYPE_CHECKING: |
28 | 38 | from bson.objectid import ObjectId |
29 | 39 | from pymongo.monitoring import _EventListeners |
30 | | - from pymongo.typings import _DocumentOut |
| 40 | + from pymongo.typings import _Address, _DocumentOut |
31 | 41 |
|
32 | 42 |
|
33 | 43 | class _CommandTelemetry: |
@@ -183,3 +193,329 @@ def failed( |
183 | 193 | service_id=self._conn.service_id, |
184 | 194 | database_name=self._dbname, |
185 | 195 | ) |
| 196 | + |
| 197 | + |
| 198 | +class _CmapTelemetry: |
| 199 | + """Combines CMAP structured logging and APM event publishing for pool and connection events.""" |
| 200 | + |
| 201 | + __slots__ = ("_address", "_client_id", "_listeners", "_should_log", "_should_publish") |
| 202 | + |
| 203 | + def __init__( |
| 204 | + self, |
| 205 | + client_id: Optional[ObjectId], |
| 206 | + address: _Address, |
| 207 | + listeners: Optional[_EventListeners], |
| 208 | + is_sdam: bool, |
| 209 | + ) -> None: |
| 210 | + self._client_id = client_id |
| 211 | + self._address = address |
| 212 | + self._listeners = listeners |
| 213 | + self._should_publish = not is_sdam and listeners is not None and listeners.enabled_for_cmap |
| 214 | + self._should_log = not is_sdam |
| 215 | + |
| 216 | + def _log(self, message: _ConnectionStatusMessage, **extra: Any) -> None: |
| 217 | + if self._should_log and _CONNECTION_LOGGER.isEnabledFor(logging.DEBUG): |
| 218 | + _debug_log( |
| 219 | + _CONNECTION_LOGGER, |
| 220 | + message=message, |
| 221 | + clientId=self._client_id, |
| 222 | + serverHost=self._address[0], |
| 223 | + serverPort=self._address[1], |
| 224 | + **extra, |
| 225 | + ) |
| 226 | + |
| 227 | + def pool_created(self, non_default_options: dict[str, Any]) -> None: |
| 228 | + # Log before publishing to prevent potential listener preemption in tests. |
| 229 | + self._log(_ConnectionStatusMessage.POOL_CREATED, **non_default_options) |
| 230 | + if self._should_publish: |
| 231 | + assert self._listeners is not None |
| 232 | + self._listeners.publish_pool_created(self._address, non_default_options) |
| 233 | + |
| 234 | + def pool_ready(self) -> None: |
| 235 | + # Log before publishing to prevent potential listener preemption in tests. |
| 236 | + self._log(_ConnectionStatusMessage.POOL_READY) |
| 237 | + if self._should_publish: |
| 238 | + assert self._listeners is not None |
| 239 | + self._listeners.publish_pool_ready(self._address) |
| 240 | + |
| 241 | + def pool_cleared(self, service_id: Optional[ObjectId], interrupt_connections: bool) -> None: |
| 242 | + # Log before publishing to prevent potential listener preemption in tests. |
| 243 | + self._log(_ConnectionStatusMessage.POOL_CLEARED, serviceId=service_id) |
| 244 | + if self._should_publish: |
| 245 | + assert self._listeners is not None |
| 246 | + self._listeners.publish_pool_cleared( |
| 247 | + self._address, |
| 248 | + service_id=service_id, |
| 249 | + interrupt_connections=interrupt_connections, |
| 250 | + ) |
| 251 | + |
| 252 | + def pool_closed(self) -> None: |
| 253 | + # Log before publishing to prevent potential listener preemption in tests. |
| 254 | + self._log(_ConnectionStatusMessage.POOL_CLOSED) |
| 255 | + if self._should_publish: |
| 256 | + assert self._listeners is not None |
| 257 | + self._listeners.publish_pool_closed(self._address) |
| 258 | + |
| 259 | + def connection_created(self, conn_id: int) -> None: |
| 260 | + # Log before publishing to prevent potential listener preemption in tests. |
| 261 | + self._log(_ConnectionStatusMessage.CONN_CREATED, driverConnectionId=conn_id) |
| 262 | + if self._should_publish: |
| 263 | + assert self._listeners is not None |
| 264 | + self._listeners.publish_connection_created(self._address, conn_id) |
| 265 | + |
| 266 | + def connection_ready(self, conn_id: int, duration: float) -> None: |
| 267 | + # Log before publishing to prevent potential listener preemption in tests. |
| 268 | + self._log( |
| 269 | + _ConnectionStatusMessage.CONN_READY, |
| 270 | + driverConnectionId=conn_id, |
| 271 | + durationMS=duration, |
| 272 | + ) |
| 273 | + if self._should_publish: |
| 274 | + assert self._listeners is not None |
| 275 | + self._listeners.publish_connection_ready(self._address, conn_id, duration) |
| 276 | + |
| 277 | + def connection_closed(self, conn_id: int, reason: str) -> None: |
| 278 | + if self._should_publish: |
| 279 | + assert self._listeners is not None |
| 280 | + self._listeners.publish_connection_closed(self._address, conn_id, reason) |
| 281 | + self._log( |
| 282 | + _ConnectionStatusMessage.CONN_CLOSED, |
| 283 | + driverConnectionId=conn_id, |
| 284 | + reason=_verbose_connection_error_reason(reason), |
| 285 | + error=reason, |
| 286 | + ) |
| 287 | + |
| 288 | + def checkout_started(self) -> None: |
| 289 | + if self._should_publish: |
| 290 | + assert self._listeners is not None |
| 291 | + self._listeners.publish_connection_check_out_started(self._address) |
| 292 | + self._log(_ConnectionStatusMessage.CHECKOUT_STARTED) |
| 293 | + |
| 294 | + def checkout_succeeded(self, conn_id: int, duration: float) -> None: |
| 295 | + if self._should_publish: |
| 296 | + assert self._listeners is not None |
| 297 | + self._listeners.publish_connection_checked_out(self._address, conn_id, duration) |
| 298 | + self._log( |
| 299 | + _ConnectionStatusMessage.CHECKOUT_SUCCEEDED, |
| 300 | + driverConnectionId=conn_id, |
| 301 | + durationMS=duration, |
| 302 | + ) |
| 303 | + |
| 304 | + def checkout_failed(self, reason: str, error: str, duration: float) -> None: |
| 305 | + if self._should_publish: |
| 306 | + assert self._listeners is not None |
| 307 | + self._listeners.publish_connection_check_out_failed(self._address, error, duration) |
| 308 | + self._log( |
| 309 | + _ConnectionStatusMessage.CHECKOUT_FAILED, |
| 310 | + reason=reason, |
| 311 | + error=error, |
| 312 | + durationMS=duration, |
| 313 | + ) |
| 314 | + |
| 315 | + def checked_in(self, conn_id: int) -> None: |
| 316 | + if self._should_publish: |
| 317 | + assert self._listeners is not None |
| 318 | + self._listeners.publish_connection_checked_in(self._address, conn_id) |
| 319 | + self._log(_ConnectionStatusMessage.CHECKEDIN, driverConnectionId=conn_id) |
| 320 | + |
| 321 | + |
| 322 | +class _HeartbeatTelemetry: |
| 323 | + """Combines SDAM structured logging and APM event publishing for server heartbeats. |
| 324 | +
|
| 325 | + The APM started event is published before connection checkout (no conn_id yet); |
| 326 | + the log entry for started is emitted after checkout once the conn_id is known. |
| 327 | + Call :meth:`apm_started` first, then :meth:`log_started` inside the checkout |
| 328 | + context, then :meth:`succeeded` or :meth:`failed` when the outcome is known. |
| 329 | + """ |
| 330 | + |
| 331 | + __slots__ = ("_address", "_awaited", "_listeners", "_publish", "_topology_id") |
| 332 | + |
| 333 | + def __init__( |
| 334 | + self, |
| 335 | + topology_id: ObjectId, |
| 336 | + address: _Address, |
| 337 | + listeners: Optional[_EventListeners], |
| 338 | + publish: bool, |
| 339 | + awaited: bool, |
| 340 | + ) -> None: |
| 341 | + self._topology_id = topology_id |
| 342 | + self._address = address |
| 343 | + self._listeners = listeners |
| 344 | + self._publish = publish |
| 345 | + self._awaited = awaited |
| 346 | + |
| 347 | + def apm_started(self) -> None: |
| 348 | + """Publish the APM heartbeat-started event (before connection checkout).""" |
| 349 | + if self._publish: |
| 350 | + assert self._listeners is not None |
| 351 | + self._listeners.publish_server_heartbeat_started(self._address, self._awaited) |
| 352 | + |
| 353 | + def log_started(self, conn_id: int, server_conn_id: Optional[int]) -> None: |
| 354 | + """Emit the log entry for heartbeat started (after connection checkout).""" |
| 355 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 356 | + _debug_log( |
| 357 | + _SDAM_LOGGER, |
| 358 | + message=_SDAMStatusMessage.HEARTBEAT_START, |
| 359 | + topologyId=self._topology_id, |
| 360 | + driverConnectionId=conn_id, |
| 361 | + serverConnectionId=server_conn_id, |
| 362 | + serverHost=self._address[0], |
| 363 | + serverPort=self._address[1], |
| 364 | + awaited=self._awaited, |
| 365 | + ) |
| 366 | + |
| 367 | + def succeeded( |
| 368 | + self, |
| 369 | + round_trip_time: float, |
| 370 | + response: Any, |
| 371 | + conn_id: int, |
| 372 | + server_conn_id: Optional[int], |
| 373 | + ) -> None: |
| 374 | + """Emit the SUCCEEDED log entry and APM event.""" |
| 375 | + if self._publish: |
| 376 | + assert self._listeners is not None |
| 377 | + self._listeners.publish_server_heartbeat_succeeded( |
| 378 | + self._address, round_trip_time, response, self._awaited |
| 379 | + ) |
| 380 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 381 | + _debug_log( |
| 382 | + _SDAM_LOGGER, |
| 383 | + message=_SDAMStatusMessage.HEARTBEAT_SUCCESS, |
| 384 | + topologyId=self._topology_id, |
| 385 | + driverConnectionId=conn_id, |
| 386 | + serverConnectionId=server_conn_id, |
| 387 | + serverHost=self._address[0], |
| 388 | + serverPort=self._address[1], |
| 389 | + awaited=self._awaited, |
| 390 | + durationMS=round_trip_time * 1000, |
| 391 | + reply=response.document, |
| 392 | + ) |
| 393 | + |
| 394 | + def failed(self, duration: float, error: Exception, conn_id: Optional[int]) -> None: |
| 395 | + """Emit the FAILED log entry and APM event.""" |
| 396 | + if self._publish: |
| 397 | + assert self._listeners is not None |
| 398 | + self._listeners.publish_server_heartbeat_failed( |
| 399 | + self._address, duration, error, self._awaited |
| 400 | + ) |
| 401 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 402 | + _debug_log( |
| 403 | + _SDAM_LOGGER, |
| 404 | + message=_SDAMStatusMessage.HEARTBEAT_FAIL, |
| 405 | + topologyId=self._topology_id, |
| 406 | + serverHost=self._address[0], |
| 407 | + serverPort=self._address[1], |
| 408 | + awaited=self._awaited, |
| 409 | + durationMS=duration * 1000, |
| 410 | + failure=error, |
| 411 | + driverConnectionId=conn_id, |
| 412 | + ) |
| 413 | + |
| 414 | + |
| 415 | +class _SdamTelemetry: |
| 416 | + """Combines SDAM structured logging and APM event publishing for topology and server events. |
| 417 | +
|
| 418 | + Topology events are queued for asynchronous delivery; log entries are emitted inline. |
| 419 | + """ |
| 420 | + |
| 421 | + __slots__ = ("_events", "_listeners", "_publish_server", "_publish_tp", "_topology_id") |
| 422 | + |
| 423 | + def __init__( |
| 424 | + self, |
| 425 | + topology_id: ObjectId, |
| 426 | + listeners: Optional[_EventListeners], |
| 427 | + events: Optional[queue.Queue[Any]], |
| 428 | + ) -> None: |
| 429 | + self._topology_id = topology_id |
| 430 | + self._listeners = listeners |
| 431 | + self._events = events |
| 432 | + self._publish_server = listeners is not None and listeners.enabled_for_server |
| 433 | + self._publish_tp = listeners is not None and listeners.enabled_for_topology |
| 434 | + |
| 435 | + def _enqueue(self, fn: Any, args: tuple[Any, ...]) -> None: |
| 436 | + if self._events is not None: |
| 437 | + self._events.put((fn, args)) |
| 438 | + |
| 439 | + def topology_opened(self) -> None: |
| 440 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 441 | + _debug_log( |
| 442 | + _SDAM_LOGGER, |
| 443 | + message=_SDAMStatusMessage.START_TOPOLOGY, |
| 444 | + topologyId=self._topology_id, |
| 445 | + ) |
| 446 | + if self._publish_tp: |
| 447 | + assert self._listeners is not None |
| 448 | + self._enqueue(self._listeners.publish_topology_opened, (self._topology_id,)) |
| 449 | + |
| 450 | + def topology_description_changed(self, old_td: Any, new_td: Any) -> None: |
| 451 | + if self._publish_tp: |
| 452 | + assert self._listeners is not None |
| 453 | + self._enqueue( |
| 454 | + self._listeners.publish_topology_description_changed, |
| 455 | + (old_td, new_td, self._topology_id), |
| 456 | + ) |
| 457 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 458 | + _debug_log( |
| 459 | + _SDAM_LOGGER, |
| 460 | + message=_SDAMStatusMessage.TOPOLOGY_CHANGE, |
| 461 | + topologyId=self._topology_id, |
| 462 | + previousDescription=repr(old_td), |
| 463 | + newDescription=repr(new_td), |
| 464 | + ) |
| 465 | + |
| 466 | + def topology_closed(self, old_td: Any, new_td: Any) -> None: |
| 467 | + """Emit APM and log events for topology description change + topology closed.""" |
| 468 | + if self._publish_tp: |
| 469 | + assert self._listeners is not None |
| 470 | + self._enqueue( |
| 471 | + self._listeners.publish_topology_description_changed, |
| 472 | + (old_td, new_td, self._topology_id), |
| 473 | + ) |
| 474 | + self._enqueue(self._listeners.publish_topology_closed, (self._topology_id,)) |
| 475 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 476 | + _debug_log( |
| 477 | + _SDAM_LOGGER, |
| 478 | + message=_SDAMStatusMessage.TOPOLOGY_CHANGE, |
| 479 | + topologyId=self._topology_id, |
| 480 | + previousDescription=repr(old_td), |
| 481 | + newDescription=repr(new_td), |
| 482 | + ) |
| 483 | + _debug_log( |
| 484 | + _SDAM_LOGGER, |
| 485 | + message=_SDAMStatusMessage.STOP_TOPOLOGY, |
| 486 | + topologyId=self._topology_id, |
| 487 | + ) |
| 488 | + |
| 489 | + def server_opened(self, address: _Address) -> None: |
| 490 | + if self._publish_server: |
| 491 | + assert self._listeners is not None |
| 492 | + self._enqueue(self._listeners.publish_server_opened, (address, self._topology_id)) |
| 493 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 494 | + _debug_log( |
| 495 | + _SDAM_LOGGER, |
| 496 | + message=_SDAMStatusMessage.START_SERVER, |
| 497 | + topologyId=self._topology_id, |
| 498 | + serverHost=address[0], |
| 499 | + serverPort=address[1], |
| 500 | + ) |
| 501 | + |
| 502 | + def server_description_changed(self, sd_old: Any, sd_new: Any, address: _Address) -> None: |
| 503 | + if self._publish_server: |
| 504 | + assert self._listeners is not None |
| 505 | + self._enqueue( |
| 506 | + self._listeners.publish_server_description_changed, |
| 507 | + (sd_old, sd_new, address, self._topology_id), |
| 508 | + ) |
| 509 | + |
| 510 | + def server_closed(self, address: _Address) -> None: |
| 511 | + if self._publish_server: |
| 512 | + assert self._listeners is not None |
| 513 | + self._enqueue(self._listeners.publish_server_closed, (address, self._topology_id)) |
| 514 | + if _SDAM_LOGGER.isEnabledFor(logging.DEBUG): |
| 515 | + _debug_log( |
| 516 | + _SDAM_LOGGER, |
| 517 | + message=_SDAMStatusMessage.STOP_SERVER, |
| 518 | + topologyId=self._topology_id, |
| 519 | + serverHost=address[0], |
| 520 | + serverPort=address[1], |
| 521 | + ) |
0 commit comments